skip to content

In RxJS 7, what replaced the deprecated multicast, publish and refCount operators, and when would you use connectable() instead of share()?

level: seniorimportance: nice to knowfreq 24%

answer

  1. three functions and a wrapper
  2. who decides when the source starts
  3. connect() returns a Subscription
  4. synchronous sources and early values
  5. resetOnDisconnect

basics

~10 s

RxJS 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 s

RxJS 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 lines
ts
import { 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

for a junior

Recall that share() and shareReplay() are the everyday sharing operators and that multicast and publish are deprecated.

for a middle

Map the deprecated multicast, publish and refCount forms to share, connectable and connect, and explain that connectable starts only on connect().

for a senior

Recognise when automatic start loses values or duplicates work, choose connectable with an owned connection, and use connect() for branching one execution.

for a principal

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