Higher-Order Observables

Preview — 3 of 10 questions

How do you create a reusable custom RxJS operator?

javascript
import { Observable, OperatorFunction, pipe } from 'rxjs';
import { filter, map, catchError, retry, tap } from 'rxjs/operators';

// Approach 1: compose existing operators with pipe()
function filterNull<T>(): OperatorFunction<T | null | undefined, T> {
  return filter((value): value is T => value != null);
}

// Approach 2: full Observable construction (for complex cases)
function retryWithDelay<T>(maxRetries: number, delayMs: number): OperatorFunction<T, T> {
  return (source: Observable<T>) => new Observable<T>(subscriber => {
    let retries = 0;
    const trySubscribe = () => {
      source.subscribe({
        next: v => subscriber.next(v),
        error: err => {
          if (retries < maxRetries) {
            retries++;
            setTimeout(trySubscribe, delayMs);
          } else {
            subscriber.error(err);
          }
        },
        complete: () => subscriber.complete(),
      });
    };
    trySubscribe();
  });
}

// Approach 3 (recommended): compose with pipe
function apiCall<T>(url: string): OperatorFunction<void, T> {
  return pipe(
    switchMap(() => this.http.get<T>(url)),
    retry({ count: 3, delay: 1000 }),
    tap(result => this.cache.set(url, result)),
    catchError(err => {
      this.logger.error(err);
      return throwError(() => err);
    }),
  );
}
AExtend the Observable class and override the pipe() method
BImplement the OperatorFunction<T, R> interface as a class
CUse the createOperator() factory function from rxjs/operators
DWrite a function that takes an Observable<T> and returns an Observable<R>, composing existing operators inside

What is a higher-order Observable and how do flattening operators handle them?

javascript
// Higher-order Observable — emits Observable<string> values
const fileUploads$: Observable<Observable<string>> = selectedFiles$.pipe(
  map(file => this.uploadService.upload(file)),  // each map emits an Observable
);

// Without flattening — you'd have to subscribe inside subscribe (callback hell):
fileUploads$.subscribe(innerObs$ => {
  innerObs$.subscribe(url => console.log(url));  // ❌ nested subscriptions
});

// With flattening operators — clean, composable:
selectedFiles$.pipe(
  mergeMap(file => this.uploadService.upload(file)),  // ✅ flattens automatically
).subscribe(url => console.log(url));

// Choosing the right flattening operator:
// mergeMap   → all inner Observables active concurrently
// switchMap  → cancels previous inner when new outer value arrives
// concatMap  → queues inner Observables, maintains order
// exhaustMap → ignores new outer values while inner is active
AAn Observable that emits values with higher priority than regular Observables
BAn Observable that emits other Observables as values, requiring a flattening operator to access the inner values
CAn Observable created from a Promise or async/await
DAn Observable that runs on a Scheduler instead of the default micro-task queue

What is an RxJS Scheduler and when would you use asyncScheduler vs queueScheduler?

javascript
import { scheduled, asyncScheduler, queueScheduler, animationFrameScheduler, asapScheduler } from 'rxjs';
import { observeOn, subscribeOn } from 'rxjs/operators';

// asyncScheduler — uses setTimeout (macro-task)
of(1, 2, 3).pipe(
  observeOn(asyncScheduler),
).subscribe(console.log);
// Values emitted in separate setTimeout calls

// queueScheduler — synchronous queue (runs before async tasks)
of(1, 2, 3).pipe(
  observeOn(queueScheduler),
).subscribe(console.log);
// Values emitted synchronously but queued (useful for recursion)

// animationFrameScheduler — uses requestAnimationFrame
// ✅ Ideal for animations/canvas updates
interval(0, animationFrameScheduler).pipe(
  map(() => this.canvas.render()),
).subscribe();

// asapScheduler — uses Promise.resolve (micro-task)
// Runs after current sync code but before macro-tasks

// subscribeOn — controls when subscription happens
// observeOn  — controls when notifications are delivered
ASchedulers define the thread on which Observables execute
BasyncScheduler is for HTTP calls; queueScheduler is for UI events
CSchedulers control the timing of Observable subscriptions and emissions; asyncScheduler dispatches via setTimeout; queueScheduler uses a synchronous queue (FIFO)
DSchedulers are deprecated in RxJS 7 — use requestAnimationFrame directly

Sign up free to play

Answer all 10 questions (7 more), see explanations for every answer, and track your score.