In RxJS 7, how do you write a reusable custom operator, both by composing existing operators with pipe() and from scratch?
answer
- a function from Observable to Observable
- the standalone pipe()
- new Observable around source.subscribe
- forward error and complete; unsubscribe source
basics
~20 sA custom operator is a function from a source Observable to a new one. Compose existing operators with the standalone pipe(); only when that cannot work, return a new Observable that forwards next, error and complete and unsubscribes the source in teardown.
solid answer
~40 sIn RxJS 7 an operator is just a function of type `OperatorFunction<T, R>` - `(source: Observable<T>) => Observable<R>` - usually returned by a factory so it can take arguments. The preferred way is composition: `export const nonNegativeCents = () => pipe(filter(v => v >= 0), map(v => Math.round(v * 100)))`, using the standalone `pipe()` exported from `'rxjs'`. When no combination of existing operators does the job, return `source => new Observable(subscriber => { ... })`: subscribe to the source, forward `next`, `error` and `complete`, keep any per-subscription state inside that function, and return a teardown that unsubscribes from the source and releases resources. `Observable.lift` is an internal detail and should not be used.
code
ts · 27 linesimport { MonoTypeOperatorFunction, Observable, OperatorFunction, filter, map, pipe } from 'rxjs';
// Option 1: composition
export function nonNegativeCents(): OperatorFunction<number, number> {
return pipe(
filter((amount) => amount >= 0),
map((amount) => Math.round(amount * 100)),
);
}
// Option 2: from scratch
export function logNumbered<T>(label: string): MonoTypeOperatorFunction<T> {
return (source) =>
new Observable<T>((subscriber) => {
let count = 0; // per subscription
const sub = source.subscribe({
next: (value) => {
count++;
console.log(label, count, value);
subscriber.next(value);
},
error: (err) => subscriber.error(err),
complete: () => subscriber.complete(),
});
return () => sub.unsubscribe();
});
}go deeper
Recall that an operator is a function from one Observable to another, and that pipe() can bundle existing operators into a new one.
Show a composed operator with the standalone pipe() and explain the OperatorFunction type and how .pipe() applies functions in order.
Write a from-scratch operator correctly: forward error and complete, keep state per subscription, unsubscribe the source in teardown, and avoid the internal lift API.
Decide which recurring pipelines deserve a shared operator library, how they are named and tested, and when a bespoke operator is less clear than inline built-ins.
## What an operator is In RxJS 7 a **pipeable operator** is an ordinary function that takes an `Observable` and returns another one. Its type is `OperatorFunction<T, R>`, or `MonoTypeOperatorFunction<T>` when input and output types match, both exported from `'rxjs'`. What people call "the `map` operator" is really a **factory**: `map(fn)` returns the operator function that `.pipe()` then calls with the source. `source$.pipe(a, b, c)` is simply `c(b(a(source$)))`. Anything with that shape can go into `.pipe()`, including your own functions. ## Option 1: compose with pipe() The RxJS guide recommends extracting a recurring sequence into a new operator with the **standalone `pipe()` function** (not the `.pipe()` method): 1. write a factory that returns `pipe(op1, op2, ...)`; 2. give it a descriptive name and explicit types; 3. use it like any built-in operator. Because it is built from existing operators, it inherits their correct handling of errors, completion and unsubscription for free. Calling `pipe()` with no arguments returns an identity function, and `.pipe()` with no arguments returns the same Observable. ## Option 2: build from scratch When no combination of existing operators expresses the behaviour - a rare case - return a function that wraps the source in `new Observable`: - **subscribe to the source inside the subscriber function**, so each subscription to the result creates its own subscription to the source; - **forward all three notifications**: `next` after your logic, and `error` and `complete` unchanged unless the operator deliberately alters them; - **return a teardown** that unsubscribes from the source and clears any timers or listeners you created; - **keep mutable state inside the subscriber function**, not in the factory, so two subscribers do not share a counter or buffer. The official guide's example re-implements `delay` this way: it tracks timers per subscription, delays `complete` until pending timers have fired, forwards `error` immediately, and clears every timer in the teardown. ## Common mistakes in hand-written operators | Mistake | Symptom | |---|---| | Not forwarding `error` | the consumer never learns of a failure; the stream appears to hang | | Not forwarding `complete` | `reduce`, `toArray` and `lastValueFrom` downstream never finish | | No teardown | the source keeps running after the consumer unsubscribes | | State declared in the factory | two subscribers corrupt each other's counters | | Using `Observable.lift` | relies on an internal API that is not part of the public surface | The RxJS 7 breaking-changes notes say `lift` was never documented for end users and recommend writing operators that return `new Observable` instead. ## Typing - Annotate the factory's return type as `OperatorFunction<In, Out>`; the compiler then checks that the composed steps line up. - Use a generic `<T>` for operators that do not care about the value type, such as a logging operator, typed `MonoTypeOperatorFunction<T>`. - Keep the factory's parameters simple values or functions, so the operator reads like the built-ins. ## Choosing between the two - If the behaviour can be described as "do A, then B" with existing operators, compose - it is shorter, and error, completion and unsubscription handling come from battle-tested code. - If the operator needs its own timing, its own buffering or custom completion rules that no combination expresses, write it from scratch - and test the error, completion and unsubscribe paths explicitly. - If you are unsure, compose first; a from-scratch rewrite can come later without changing callers, since both forms have the same function type. ## Imports in RxJS 7 Since RxJS 7.2 every operator is exported from `'rxjs'` itself, and the RxJS guide describes the `'rxjs/operators'` entry point as deprecated. New code imports `map`, `filter`, `pipe` and `Observable` from `'rxjs'`. ## Testing a custom operator A custom operator is a pure function of its source, so it tests well: feed it a known source such as `of(...)`, collect the output, and assert on values and completion. Timing-dependent operators are best tested with virtual time, which is its own topic. ## Interview summary Compose with `pipe()` first; write from scratch only when necessary, and then forward all notifications, keep state per subscription, and clean up in the teardown.
- In an RxJS custom operator written with new Observable, why must the counter be declared inside the subscriber function?The factory and the returned operator function run once when the pipe is built, but the subscriber function runs once per subscription. A counter declared outside it is shared by every subscriber, so two subscriptions would interleave their counts. Declaring it inside gives each subscription its own state.
- In RxJS, what is the difference between the standalone pipe() function and the Observable .pipe() method?`source$.pipe(a, b)` applies the operators to that Observable and returns the resulting Observable. The standalone `pipe(a, b)`, imported from `'rxjs'`, applies nothing yet: it returns a new function that will apply `a` then `b` to whatever source it is given later, which is exactly the shape of a reusable operator.
saying these in an interview costs you the question
- Custom operators must extend an RxJS Operator class and override call().
- Observable.lift is the recommended public API for custom operators in RxJS 7.
- A from-scratch operator only needs to forward next; error and complete propagate automatically.
- State declared in the operator factory is automatically separate for each subscriber.
- Operators must still be imported from 'rxjs/operators' in RxJS 7.