Создание Observable потоков

Angular активно использует RxJS, предоставляя возможности для работы с асинхронными потоками данных через Observable. Observable представляет собой последовательность событий, которые могут приходить асинхронно, и позволяет их подписывать, трансформировать и комбинировать.

Основные способы создания Observable

  1. Конструктор Observable
import { Observable } from 'rxjs';

const observable = new Observable<number>(subscriber => {
  subscriber.next(1);
  subscriber.next(2);
  subscriber.complete();
});
  • subscriber.next(value) – отправка нового значения.
  • subscriber.error(error) – уведомление об ошибке.
  • subscriber.complete() – завершение потока.
  1. Операторы создания
  • of(...values) – создает Observable из переданных значений.
  • from(array | promise | iterable) – превращает массив, промис или итерируемый объект в поток.
  • interval(ms) – эмитирует числа с заданным интервалом времени.
  • timer(dueTime, period) – запускает событие через dueTime и затем с периодом period.

Пример:

import { of, from, interval } from 'rxjs';

const numbers$ = of(1, 2, 3);
const array$ = from([10, 20, 30]);
const timer$ = interval(1000);

Подписка на Observable

Для обработки значений используется метод subscribe:

numbers$.subscribe({
  next: value => console.log('Получено:', value),
  error: err => console.error('Ошибка:', err),
  complete: () => console.log('Поток завершен')
});

Подписка позволяет реагировать на каждый элемент потока, обрабатывать ошибки и фиксировать завершение последовательности.

Операторы трансформации и комбинирования

RxJS предоставляет более сотни операторов, которые делят на категории:

  • Трансформация: map, filter, scan.
  • Комбинирование: merge, concat, combineLatest.
  • Потоковое управление временем: debounceTime, throttleTime, delay.

Пример использования map и filter:

import { from } from 'rxjs';
import { map, filter } from 'rxjs/operators';

from([1, 2, 3, 4, 5]).pipe(
  filter(x => x % 2 === 0),
  map(x => x * 10)
).subscribe(console.log);
// Вывод: 20, 40

Управление подписками

Для предотвращения утечек памяти важно отписываться от Observable, особенно в Angular компонентах:

import { Subscription } from 'rxjs';

export class ExampleComponent implements OnDestroy {
  private subscription: Subscription;

  ngOnInit() {
    this.subscription = interval(1000).subscribe(console.log);
  }

  ngOnDestroy() {
    this.subscription.unsubscribe();
  }
}

Angular также поддерживает асинхронный пайп async, который автоматически управляет подпиской и отпиской:

<div *ngFor="let item of items$ | async">{{ item }}</div>

Создание сложных потоков

Observable позволяют строить сложные асинхронные сценарии:

  • последовательные HTTP-запросы с switchMap,
  • комбинирование нескольких потоков с forkJoin или combineLatest,
  • реализация debounce и throttle для обработки ввода пользователя.

Использование Observable в Angular обеспечивает реактивный подход к разработке, позволяя управлять асинхронными данными, событиями и состояниями компонентов гибко и эффективно.