Instruction file imported from wghglory/ngx-lift-workspace (
.cursor/rules/rxjs.mdc). Copyright stays with the author.
RxJS Patterns and Best Practices
Observable Creation
-
Prefer Signals: Use Angular Signals for component state when possible
-
Observables for Streams: Use Observables for continuous data streams (HTTP, events, timers)
-
Creation Functions: Use appropriate creation functions
// ✅ Good data$ = of(data); data$ = from(promise); data$ = fromEvent(element, 'click'); // ❌ Avoid data$ = new Observable(subscriber => { ... });
Operators
Common Operators
-
map: Transform values
users$ = this.http .get<User[]>(url) .pipe(map((users) => users.map((u) => ({...u, fullName: `${u.firstName} ${u.lastName}`})))); -
filter: Filter values
activeUsers$ = users$.pipe(filter((user) => user.active)); -
switchMap: Cancel previous inner observables (HTTP requests)
user$ = userId$.pipe(switchMap((id) => this.http.get<User>(`/users/${id}`))); -
mergeMap/flatMap: Concurrent inner observables
results$ = ids$.pipe(mergeMap((id) => this.http.get<Result>(`/results/${id}`))); -
concatMap: Sequential inner observables
results$ = ids$.pipe(concatMap((id) => this.http.get<Result>(`/results/${id}`))); -
exhaustMap: Ignore new emissions while inner observable is active
save$ = saveAction$.pipe(exhaustMap(() => this.http.post('/save', data))); -
catchError: Handle errors
data$ = this.http.get(url).pipe( catchError((error) => { console.error(error); return of(defaultValue); }), ); -
tap: Side effects (logging, debugging)
data$ = this.http.get(url).pipe( tap((data) => console.log('Received:', data)), tap({error: (err) => console.error('Error:', err)}), );
ngx-lift Operators
-
createAsyncState: Transform Observable to AsyncState
usersState$ = this.userService.getUsers().pipe( createAsyncState({ next: (users) => console.log('Loaded users:', users), error: (err) => console.error('Failed to load users:', err), }), ); -
switchMapWithAsyncState: Combine switchMap with async state
userState$ = userId$.pipe(switchMapWithAsyncState((id) => this.userService.getUser(id))); -
combineLatestEager: Combine observables with initial values
vm$ = combineLatestEager({ users: this.users$, filters: this.filters$, }); -
poll: Polling with configurable interval
dataState$ = poll({ interval: 5000, pollingFn: () => this.http.get('/data'), forceRefresh: this.refresh$, initialValue: {isLoading: false, error: null, data: null}, }); -
distinctOnChange: Execute callback on value change
data$.pipe( distinctOnChange((prev, curr) => { console.log('Value changed from', prev, 'to', curr); }), );
Async State Management
-
AsyncState Pattern: Use AsyncState for loading/error/success states
interface AsyncState<T, E = Error> { status: ResourceStatus; isLoading: boolean; error: E | null; data: T | null; } -
Template Usage: Use async pipe with AsyncState
@if (userState$ | async; as state) { @if (state.isLoading) { <cll-spinner /> } @if (state.error) { <cll-alert [error]="state.error" /> } @if (state.data; as user) { <p>{{ user.name }}</p> } }
Subscription Management
-
Async Pipe: Prefer async pipe in templates (automatic subscription management)
@if (data$ | async; as data) { <p>{{ data.name }}</p> } -
Manual Subscriptions: Only when necessary, always unsubscribe
private destroyRef = inject(DestroyRef); ngOnInit() { const subscription = this.data$.subscribe(data => { // handle data }); this.destroyRef.onDestroy(() => { subscription.unsubscribe(); }); } -
takeUntil Pattern: Use
takeUntilwithDestroyReforSubjectprivate destroy$ = new Subject<void>(); ngOnInit() { this.data$.pipe( takeUntil(this.destroy$) ).subscribe(); } ngOnDestroy() { this.destroy$.next(); this.destroy$.complete(); }
Subjects
-
BehaviorSubject: For state that needs initial value
private userSubject = new BehaviorSubject<User | null>(null); user$ = this.userSubject.asObservable(); -
Subject: For events/actions
private saveAction$ = new Subject<void>(); save() { this.saveAction$.next(); } -
ReplaySubject: For caching last N values
private dataSubject = new ReplaySubject<Data>(3);
Error Handling
-
catchError: Always handle errors in HTTP requests
data$ = this.http.get(url).pipe( catchError((error) => { this.errorService.handle(error); return of(null); }), ); -
retry: Retry failed operations
data$ = this.http.get(url).pipe( retry(3), catchError((error) => of(null)), ); -
retryWhen: Retry with custom logic
data$ = this.http.get(url).pipe(retryWhen((errors) => errors.pipe(delay(1000), take(3))));
Combining Observables
-
combineLatest: Combine latest values from multiple observables
combined$ = combineLatest([users$, filters$]).pipe(map(([users, filters]) => filterUsers(users, filters))); -
merge: Merge multiple observables
merged$ = merge(source1$, source2$); -
zip: Combine observables in order
zipped$ = zip(users$, roles$).pipe(map(([users, roles]) => combineUsersAndRoles(users, roles))); -
forkJoin: Wait for all observables to complete
data$ = forkJoin({ users: this.getUsers(), roles: this.getRoles(), });
Signals and Observables
-
toSignal: Convert Observable to Signal
user = toSignal(this.user$); user = toSignal(this.user$, {initialValue: null}); -
toObservable: Convert Signal to Observable
user$ = toObservable(this.user); -
combineFrom: Combine signals and observables (ngx-lift)
combined = combineFrom([this.signal, this.observable$]); combined = combineFrom({a: this.signal, b: this.observable$});
Performance Optimization
-
shareReplay: Share and replay values
data$ = this.http.get(url).pipe(shareReplay(1)); -
distinctUntilChanged: Emit only when value changes
filtered$ = data$.pipe(distinctUntilChanged((prev, curr) => prev.id === curr.id)); -
debounceTime: Debounce rapid emissions
searchResults$ = searchTerm$.pipe( debounceTime(300), switchMap((term) => this.search(term)), ); -
throttleTime: Throttle emissions
scrollEvents$ = fromEvent(window, 'scroll').pipe(throttleTime(100));
Best Practices
- Avoid Manual Subscriptions: Use async pipe or reactive patterns
- Error Handling: Always handle errors in observable chains
- Memory Leaks: Always unsubscribe from long-lived subscriptions
- Operator Order: Order operators logically (transformation → filtering → side effects)
- Type Safety: Maintain type safety throughout observable chains
- Documentation: Document complex observable chains
