Observable и Observer паттерны

Angular активно использует реактивное программирование через библиотеку RxJS, где ключевыми элементами являются паттерны Observable и Observer.

Observable

Observable — это поток данных, который может быть асинхронным или синхронным. Он описывает, как данные будут доставляться со временем. Observable может передавать три вида уведомлений:

  • next(value) — новое значение потока.
  • error(err) — уведомление об ошибке, прерывающее поток.
  • complete() — уведомление о завершении передачи данных.

Создание Observable в Angular:

import { Observable } from 'rxjs';

const data$ = new Observable<number>(observer => {
  observer.next(1);
  observer.next(2);
  observer.complete();
});

Observer

Observer — объект, подписывающийся на Observable, чтобы получать уведомления о новых данных, ошибках или завершении. Содержит методы: next, error, complete.

Подписка на Observable:

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

data$.subscribe(observer);

Операторы и управление потоками

RxJS предоставляет операторы, позволяющие трансформировать и фильтровать данные:

  • map() — преобразование значений потока.
  • filter() — фильтрация данных по условию.
  • switchMap() — переключение на новый поток при каждом новом значении.
  • merge() — объединение нескольких Observable.

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

import { Injectable } from '@angular/core';
import { HttpClient } from '@angular/common/http';
import { Observable } from 'rxjs';
import { map } from 'rxjs/operators';

@Injectable({ providedIn: 'root' })
export class DataService {
  constructor(private http: HttpClient) {}

  getTransformedData(): Observable<string[]> {
    return this.http.get<{ name: string }[]>('/api/items').pipe(
      map(items => items.map(item => item.name.toUpperCase()))
    );
  }
}

Особенности работы с Observable

  • Ленивость: Observable не выполняются, пока на них не подписались.
  • Многоразовость подписок: каждый подписчик получает отдельный поток данных. Для совместного использования данных применяется оператор shareReplay().
  • Управление памятью: подписки необходимо отписывать при уничтожении компонента (ngOnDestroy) для предотвращения утечек памяти.

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