In RxJS 7, what replaced the deprecated multicast, publish and refCount operators, and when would you use connectable() instead of share()?
answer
- three functions and a wrapper
- who decides when the source starts
- connect() returns a Subscription
- synchronous sources and early values
- resetOnDisconnect
basics
~10 sRxJS 7 reduced multicasting to connectable, connect, share and shareReplay; use connectable() when you must attach every subscriber before the source starts, because share() starts it on the first subscription.
solid answer
~30 sRxJS 7 deprecated `ConnectableObservable`, `multicast`, `publish`, `publishBehavior`, `publishLast`, `publishReplay` and `refCount`; they are slated for removal in v8. The replacements: `multicast(...)` plus `refCount()` becomes `share({ connector })`; a connectable observable becomes `connectable(source$, { connector })`; `multicast` with a selector becomes the `connect(selector)` operator. `share()` starts the source **automatically** on the first subscription, so with a source that emits synchronously the first subscriber can receive everything before a second one attaches. `connectable()` starts the source only when you call its `connect()` method, which returns the `Subscription` you use to disconnect. That explicit start is the reason to choose it.
code
ts · 13 linesimport { Observable, connect, filter, map, merge } from 'rxjs';
interface Message { room: string; text: string; mention: boolean }
declare const messages$: Observable<Message>;
const view$ = messages$.pipe(
connect((shared$) =>
merge(
shared$.pipe(map((m) => ({ kind: 'line' as const, m }))),
shared$.pipe(filter((m) => m.mention), map((m) => ({ kind: 'ping' as const, m }))),
),
),
);go deeper
Recall that share() and shareReplay() are the everyday sharing operators and that multicast and publish are deprecated.
Map the deprecated multicast, publish and refCount forms to share, connectable and connect, and explain that connectable starts only on connect().
Recognise when automatic start loses values or duplicates work, choose connectable with an owned connection, and use connect() for branching one execution.
Weigh explicit connection management against automatic reference counting across a codebase, and plan removing RxJS 6 multicasting before a v8 upgrade.
## The RxJS 7 multicasting API RxJS 6 had many overlapping ways to multicast. RxJS 7 reduced them to three functions plus one wrapper: - **`share(config?)`** — an operator that multicasts automatically, with reference counting and configurable resets. - **`shareReplay(...)`** — a thin wrapper over `share` that uses a `ReplaySubject`. - **`connectable(source, config?)`** — a function that returns an observable you start **by hand** with `connect()`. - **`connect(selector, config?)`** — an operator that multicasts the source **inside** a selector function, for composing several branches of one execution. Everything else — `ConnectableObservable`, `multicast`, `publish`, `publishBehavior`, `publishLast`, `publishReplay` and `refCount` — was deprecated in RxJS 7.0 and is slated for removal in RxJS 8. ## Mapping old code to new | Deprecated form | Replacement | |---|---| | `multicast(() => new Subject()), refCount()` | `share({ connector: () => new Subject() })` | | `multicast(() => new Subject())` used as a connectable | `connectable(source$, { connector: () => new Subject() })` | | `publish()` used as a connectable | `connectable(source$, { connector: () => new Subject(), resetOnDisconnect: false })` | | `new ConnectableObservable(source$, factory)` | `connectable(source$, { connector: factory })` | | `multicast(subject, selector)` / `publish(selector)` | `connect(selector)` | ## share versus connectable: who starts the source **`share()`** starts the source the moment the **first** subscriber arrives, and by default stops it when the **last** one leaves. That is ideal for streams consumers come and go from. It has one sharp edge: with a source that emits **synchronously**, the first subscriber can receive every value, and the source can complete, before the next line of code subscribes a second consumer. Because `share()` resets on completion by default, the second subscriber then triggers a **second** execution. **`connectable(source$)`** separates subscribing from starting: 1. Subscribing to the connectable observable only attaches you to its internal subject; the source does not run. 2. Calling `connect()` subscribes the subject to the source. It returns a `Subscription`. 3. Unsubscribing that `Subscription` disconnects the source. 4. With the default `resetOnDisconnect: true`, a fresh subject is created on disconnection so it can be connected again; with `false`, the same subject stays in place. ```ts import { connectable, of, tap } from 'rxjs'; const source$ = of(1, 2, 3).pipe(tap(() => console.log('produced'))); const shared$ = connectable(source$); shared$.subscribe((v) => console.log('A', v)); shared$.subscribe((v) => console.log('B', v)); shared$.connect(); // produced, A 1, B 1, produced, A 2, B 2, produced, A 3, B 3 ``` Both subscribers see all three values from a single execution — which `share()` could not guarantee for this synchronous source. ## The connect operator `connect(selector)` multicasts the source **only for the duration of the selector's result**. The selector receives a shared observable and returns a combined observable; the source is subscribed once, and all branches see the same values: - It keeps multicasting local to one pipeline, with no connection to manage. - It is the replacement for the selector overloads of `multicast` and `publish`. - A typical use is splitting one expensive stream into several filtered branches and merging them back. ## Choosing - **Consumers come and go; start and stop automatically:** `share()` (or `shareReplay` for replay). - **All consumers must be attached before the first value, or the start time is a deliberate decision:** `connectable()` and an explicit `connect()`. - **Several branches of one execution within a single pipeline:** `connect(selector)`. With `connectable`, you own the lifetime: keep the `Subscription` that `connect()` returned and unsubscribe it, or the source runs until it completes on its own.
- What does connect() on a connectable observable return, and why keep it?It returns the `Subscription` of the source to the internal subject. Unsubscribing it disconnects the source; with the default `resetOnDisconnect: true` a fresh subject is created so it can be connected again. If you drop that `Subscription`, nothing ever stops the source short of its own completion.
- Why can share() run a synchronous source twice for two subscribers written on consecutive lines?`share()` subscribes to the source when the first subscriber arrives. A synchronous source such as `of(1, 2, 3)` emits everything and completes inside that first `subscribe()` call. `share()` resets on completion by default, so the second `subscribe()` starts a brand-new execution. `connectable()` avoids this by starting only on `connect()`, after both are attached.
saying these in an interview costs you the question
- publish() and refCount() are still the recommended way to multicast in RxJS 7
- connectable() starts the source as soon as the first subscriber arrives
- share() guarantees that every subscriber sees every value of a synchronous source
- Unsubscribing every subscriber of a connectable observable stops the source
- The connect operator and connectable() are two names for the same thing