Pular para o conteúdo
< samuelsantana.dev />
Voltar para o BlogDiagrama de marbles do typeahead: sete teclas de 'a' a 'angular', seis requisições canceladas pelo switchMap e só a resposta 'angular@520ms' chegando à tela.

Performance com RxJS, medida: o gargalo quase nunca é o operador

Samuel Santana
Publicado em 05 de outubro de 2026
AngularPerformanceTypeScript

Uma pergunta que aparece sob "RxJS performance" é se o RxJS é lento. Eu medi. Na minha máquina, cada operador de uma cadeia custa entre 8 e 12 nanossegundos por valor. Não é de graça, mas quase nunca é o problema. O que eu vejo custar caro numa aplicação Angular é outra coisa: o mesmo request disparado três vezes, uma fonte que continua rodando depois que a tela foi embora, um valor derivado que emite duas vezes (uma delas errada) e uma busca que termina mostrando o resultado de uma digitação antiga.

Este post mede cada um desses casos com um script curto que dá para rodar. A versão é o RxJS 7.8.2, a latest no npm hoje, 05/10/2026. Existe um 9.0.0-beta.0 na tag next desde 04/08/2026, que eu não medi; e o Angular 22.2.1, a versão atual, declara o rxjs como peer em ^6.5.3 || ^7.4.0, então a 9 ainda não é opção num projeto Angular. Rodei tudo no Node 24.19.0, num Intel Core i9-13900H com Windows. Os tempos dependem da máquina. As contagens (quantos requests, quantas emissões, quantos teardowns) não dependem: são semântica do RxJS e dão o mesmo resultado em qualquer lugar.

Ele continua dois posts anteriores: RxJS vs. Signals, de julho, e os operators do rx-state-bridge, de agosto.

Primeiro, a régua: quanto custa um operador

Antes de falar do que custa caro, vale medir o que todo mundo suspeita. O script empurra um milhão de valores por cadeias de 0, 5 e 10 operadores, alternando map e filter (o filtro sempre deixa passar, então todo valor percorre a cadeia inteira). Depois compara com as mesmas dez etapas escritas como funções comuns, e mede montar e desmontar uma assinatura numa cadeia de cinco operadores, que é o que um async pipe faz em cada linha de uma lista. Cada caso tem cinco rodadas de aquecimento e 21 rodadas medidas; a saída mostra a mediana, o mínimo e o máximo.

// node --expose-gc cost.mjs
import { Subject, filter, map } from 'rxjs';

function bench(label, iterations, fn, rounds = 21) {
  for (let i = 0; i < 5; i++) fn(iterations); // aquecimento
  const ns = [];
  for (let i = 0; i < rounds; i++) {
    const t0 = process.hrtime.bigint();
    fn(iterations);
    ns.push(Number(process.hrtime.bigint() - t0) / iterations);
  }
  ns.sort((a, b) => a - b);
  console.log(`${label.padEnd(38)} median ${ns[rounds >> 1].toFixed(1)} ns (min ${ns[0].toFixed(1)}, max ${ns.at(-1).toFixed(1)})`);
}

// map e filter alternados; o filter sempre deixa passar, então todo valor percorre a cadeia toda.
const ops = (n) => Array.from({ length: n }, (_, i) => (i % 2 ? filter((x) => x >= 0) : map((x) => x + 1)));
let sink = 0;

for (const n of [0, 5, 10]) {
  const subject = new Subject();
  subject.pipe(...ops(n)).subscribe((x) => (sink += x));
  bench(`emit through ${n} operators`, 1e6, (it) => { for (let i = 0; i < it; i++) subject.next(i); });
}

const fns = Array.from({ length: 10 }, (_, i) => (i % 2 ? (x) => (x >= 0 ? x : x) : (x) => x + 1));
bench('same 10 steps, plain functions', 1e6, (it) => {
  for (let i = 0; i < it; i++) { let x = i; for (const f of fns) x = f(x); sink += x; }
});

const piped$ = new Subject().pipe(...ops(5));
bench('subscribe + unsubscribe, 5 operators', 1e5, (it) => {
  for (let i = 0; i < it; i++) piped$.subscribe((x) => (sink += x)).unsubscribe();
});

global.gc();
const before = process.memoryUsage().heapUsed;
const subs = Array.from({ length: 1e5 }, () => piped$.subscribe((x) => (sink += x)));
global.gc();
console.log(`heap per live subscription, 5 operators: ${((process.memoryUsage().heapUsed - before) / 1e5).toFixed(0)} bytes`);
subs.forEach((s) => s.unsubscribe());
emit through 0 operators               median 22.7 ns (min 21.0, max 25.4)
emit through 5 operators               median 64.4 ns (min 62.5, max 66.5)
emit through 10 operators              median 120.7 ns (min 107.3, max 150.2)
same 10 steps, plain functions         median 30.9 ns (min 29.1, max 33.8)
subscribe + unsubscribe, 5 operators   median 494.6 ns (min 475.1, max 595.7)
heap per live subscription, 5 operators: 3265 bytes

Rodei o benchmark cinco vezes (uma com este script, quatro com uma versão anterior que só difere nos rótulos). A mediana dos 10 operadores ficou entre 121 e 144 ns; a de montar e desmontar, entre 484 e 740 ns. A memória por assinatura deu 3.265 ou 3.266 bytes em todas.

Traduzindo os números:

  • Cada operador acrescenta de 8 a 12 ns por valor. Dez operadores custam cerca de quatro vezes as mesmas etapas como funções diretas. O overhead existe.
  • A 60 Hz, um quadro dura 16,7 ms. Mesmo na pior mediana (144 ns por valor), a cadeia de dez operadores precisaria de mais de 110 mil emissões dentro de um único quadro para ocupá-lo inteiro sozinha.
  • Uma lista de mil linhas com um async pipe por linha, numa cadeia de cinco operadores, leva entre 0,5 e 0,75 ms para montar as mil assinaturas e ocupa cerca de 3 MB de heap (mil vezes 3.265 bytes) enquanto elas existem.

O que o número não diz: ele mede só o RxJS, sem a detecção de mudanças do Angular e sem DOM. E foi medido no Node. O Node usa o V8, o mesmo motor do Chrome, então a ordem de grandeza deve se manter no Chrome; o valor absoluto muda com a CPU e com o estado do JIT. Firefox e Safari usam outros motores, e eu não medi neles.

Isso dá escala a uma frase do post de julho, que conta o overhead de memória de cada async pipe entre as dores do RxJS: "em listas longas, esse custo se acumula". Acumula, mas devagar: mil assinaturas custam menos de um milissegundo para montar. O que pesa numa lista longa é o trabalho que cada assinatura dispara. O resto do post é sobre isso.

O mesmo request, três vezes

O HttpClient devolve Observables frios. A documentação do Angular é direta: nenhuma requisição acontece até alguém assinar, e assinar o mesmo Observable várias vezes dispara várias requisições ao backend, uma por assinatura. Um template com user$ | async em três lugares faz três requests.

O script abaixo imita esse comportamento com um Observable que conta quantas vezes foi executado. Três consumidores assinam juntos, como três async pipes no mesmo template; um quarto assina depois que a resposta chegou, como um componente filho criado mais tarde.

import { Observable, share, shareReplay } from 'rxjs';

const sleep = (ms) => new Promise((r) => setTimeout(r, ms));

// Igual ao HttpClient: nada acontece até o subscribe, e cada subscribe é um request.
const fakeGet = (stats) =>
  new Observable((subscriber) => {
    stats.requests++;
    const t = setTimeout(() => { subscriber.next({ name: 'Ada' }); subscriber.complete(); }, 50);
    return () => clearTimeout(t);
  });

async function scenario(label, operator) {
  const stats = { requests: 0 };
  const user$ = operator ? fakeGet(stats).pipe(operator) : fakeGet(stats);
  for (let i = 0; i < 3; i++) user$.subscribe(); // três `user$ | async` no mesmo template
  await sleep(100);
  const atOnce = stats.requests;
  user$.subscribe(); // um componente criado depois que a resposta chegou
  await sleep(100);
  console.log(`${label.padEnd(46)} 3 at once: ${atOnce} | +1 late: ${stats.requests} total`);
}

await scenario('no sharing', null);
await scenario('share()', share());
await scenario('shareReplay(1)', shareReplay(1));
await scenario('shareReplay({ bufferSize: 1, refCount: true })', shareReplay({ bufferSize: 1, refCount: true }));
no sharing                                     3 at once: 3 | +1 late: 4 total
share()                                        3 at once: 1 | +1 late: 2 total
shareReplay(1)                                 3 at once: 1 | +1 late: 1 total
shareReplay({ bufferSize: 1, refCount: true }) 3 at once: 1 | +1 late: 1 total

Três coisas nessa tabela merecem atenção.

share() resolve quem chega junto, mas não quem chega depois. O quarto assinante disparou outro request. É o padrão do operador: no código do share, resetOnComplete começa como true, e a documentação de ShareConfig explica que, com ele ligado, o Observable volta ao estado "frio" quando a fonte completa. Um request HTTP completa logo depois da resposta.

shareReplay(1) guarda o último valor e entrega a quem chega depois. A documentação do shareReplay avisa o preço: uma fonte que completou com sucesso fica em cache no Observable compartilhado para sempre. Uma fonte que terminou em erro pode ser tentada de novo.

refCount: true não desliga esse cache. O quarto assinante recebeu o valor guardado sem request novo, mesmo com refCount: true e mesmo com a contagem de assinantes tendo chegado a zero antes dele. O motivo está no código do share: o reset por contagem zero só acontece se a fonte ainda não completou nem falhou. Ou seja, refCount decide o que acontece com uma fonte viva quando todo mundo sai. Ele não decide se o resultado de um request fica guardado.

No template, a correção mais barata nem usa operador. É uma assinatura só, com alias:

@if (user$ | async; as user) {
  <app-avatar [user]="user" />
  <h2>{{ user.name }}</h2>
  <app-permissions [roles]="user.roles" />
}

A outra é chamar toSignal uma vez, num campo da classe, e ler o signal onde precisar. A documentação de interoperabilidade avisa que toSignal cria uma assinatura, então não deve ser chamado repetidamente para o mesmo Observable. Um detalhe do alias: o @if não renderiza o bloco quando o valor é falsy. Para um objeto de usuário, tudo bem; para um contador que pode valer zero, não.

shareReplay sem refCount: vazamento, não lentidão

O refCount padrão do shareReplay é false. A documentação diz o que isso significa: quando a contagem de assinantes chega a zero, a fonte não é desassinada, e o ReplaySubject interno pode rodar para sempre. É uma escolha deliberada, pensada para manter vivas fontes caras de montar. O problema aparece quando ela é usada onde ninguém queria isso.

Medi dois casos. No caso A, o usuário sai da tela 20 ms depois de o request começar, e o request levaria 100 ms. No caso B, a fonte nunca completa: um interval faz o papel de polling ou websocket. A tela é montada e desmontada mil vezes, e em cada montagem dois assinantes entram e saem.

// node --expose-gc refcount.mjs
import { Observable, interval, shareReplay } from 'rxjs';

const sleep = (ms) => new Promise((r) => setTimeout(r, ms));
const timers = () => process.getActiveResourcesInfo().filter((r) => r === 'Timeout').length;
const heapMiB = () => (global.gc(), process.memoryUsage().heapUsed / 2 ** 20);

const variants = {
  'shareReplay(1)': () => shareReplay(1),
  'shareReplay({ bufferSize: 1, refCount: true })': () => shareReplay({ bufferSize: 1, refCount: true }),
};

// A) O usuário sai enquanto o request ainda está em andamento.
for (const [label, op] of Object.entries(variants)) {
  let aborted = 0, processed = 0;
  const req$ = new Observable((s) => {
    let done = false;
    const t = setTimeout(() => { done = true; processed++; s.next('data'); s.complete(); }, 100);
    return () => { if (!done) aborted++; clearTimeout(t); };
  }).pipe(op());
  const sub = req$.subscribe();
  await sleep(20);
  sub.unsubscribe(); // componente destruído
  await sleep(150);
  console.log(`A ${label.padEnd(46)} aborted: ${aborted} | processed anyway: ${processed}`);
}

// B) Uma fonte que nunca completa (polling, websocket), mil ciclos de montar e desmontar.
for (const [label, op] of Object.entries(variants)) {
  const baseTimers = timers(), baseHeap = heapMiB();
  let ticks = 0, teardowns = 0;
  for (let i = 0; i < 1000; i++) {
    const ticker$ = new Observable((s) => {
      const sub = interval(10).subscribe((n) => { ticks++; s.next(n); });
      return () => { teardowns++; sub.unsubscribe(); };
    }).pipe(op());
    const a = ticker$.subscribe(), b = ticker$.subscribe();
    a.unsubscribe(); b.unsubscribe(); // todos os consumidores saíram
  }
  const atExit = ticks;
  await sleep(500);
  console.log(`B ${label.padEnd(46)} teardowns: ${teardowns}/1000 | live timers: ${timers() - baseTimers}` +
    ` | callbacks in the next 500 ms: ${ticks - atExit} | retained: ${(heapMiB() - baseHeap).toFixed(2)} MiB`);
}
process.exit(0); // os intervals vazados manteriam o Node vivo para sempre
A shareReplay(1)                                 aborted: 0 | processed anyway: 1
A shareReplay({ bufferSize: 1, refCount: true }) aborted: 1 | processed anyway: 0
B shareReplay(1)                                 teardowns: 0/1000 | live timers: 1000 | callbacks in the next 500 ms: 32000 | retained: 2.66 MiB
B shareReplay({ bufferSize: 1, refCount: true }) teardowns: 1000/1000 | live timers: 0 | callbacks in the next 500 ms: 0 | retained: 0.05 MiB

No caso A, com shareReplay(1), o request foi até o fim e a resposta foi processada para ninguém. Com refCount: true, ele foi abortado. No HttpClient, desassinar antes da resposta aborta a requisição em andamento (está na mesma documentação de requisições), mas só se o unsubscribe chegar até a fonte. O shareReplay(1) segura esse unsubscribe no meio do caminho.

No caso B, o teardown da fonte não rodou nenhuma vez em mil montagens. Ficaram mil timers vivos, que dispararam 32 mil callbacks nos 500 ms seguintes (nesta máquina Windows, o intervalo de 10 ms disparou a cada 16 ms, mais ou menos), e 2,66 MiB de heap continuaram retidos depois da coleta de lixo, contra 0,05 MiB com refCount: true. Isso não é lentidão de operador. É trabalho que nunca para, e que cresce a cada navegação.

Note que, para memória, o caso A é quase inofensivo: o request completa e a fonte para sozinha. O vazamento de verdade precisa de uma fonte que não completa. Polling, websocket, valueChanges de formulário, Router.events e seletores de store são fontes assim. A regra que eu sigo: em componente, ou em qualquer coisa que viva menos que a aplicação, shareReplay({ bufferSize: 1, refCount: true }). refCount: false só para um cache que deve durar a aplicação inteira, e de propósito.

takeUntilDestroyed no lugar errado

O takeUntilDestroyed completa o Observable quando o contexto que o chamou (componente, diretiva, serviço) é destruído. No código do @angular/core 22.2.1, a implementação é literalmente source.pipe(takeUntil(destroyed$)). Por isso a posição dele no pipe importa tanto quanto a do takeUntil, e um Subject fazendo o papel do DestroyRef reproduz o comportamento:

import { Subject, interval, switchMap, takeUntil } from 'rxjs';

const sleep = (ms) => new Promise((r) => setTimeout(r, ms));
const timers = () => process.getActiveResourcesInfo().filter((r) => r === 'Timeout').length;

async function scenario(label, build) {
  const route$ = new Subject();     // por exemplo, parâmetros de rota
  const destroyed$ = new Subject(); // por exemplo, DestroyRef.onDestroy
  const pollOrder = () => interval(10); // consulta um endpoint a cada 10 ms
  const baseTimers = timers();
  let received = 0;
  build(route$, destroyed$, pollOrder).subscribe(() => received++);

  route$.next('order-42');
  await sleep(100);
  destroyed$.next(); // componente destruído
  const atDestroy = received;
  await sleep(300);
  console.log(`${label.padEnd(34)} values after destroy: ${received - atDestroy} | timers still running: ${timers() - baseTimers}`);
}

await scenario('takeUntil before switchMap', (route$, destroyed$, poll) =>
  route$.pipe(takeUntil(destroyed$), switchMap(poll)));
await scenario('takeUntil after switchMap (last)', (route$, destroyed$, poll) =>
  route$.pipe(switchMap(poll), takeUntil(destroyed$)));
process.exit(0); // o interval vazado manteria o processo vivo
takeUntil before switchMap         values after destroy: 20 | timers still running: 1
takeUntil after switchMap (last)   values after destroy: 0 | timers still running: 0

Com o takeUntil antes do switchMap, a destruição completou o fluxo externo, mas o polling interno continuou: 20 valores em 300 ms depois do destroy, e o timer ainda rodando quando o script terminou. O motivo está no código do switchMap: a saída só completa quando o fluxo externo completou e não há assinatura interna ativa. Um polling nunca completa, então a saída também não. Com o takeUntil por último, o unsubscribe desce pela cadeia inteira e o polling para.

A regra, então: takeUntilDestroyed() é o último operador antes do subscribe. Ou não assine à mão e deixe o async pipe ou o toSignal cuidarem disso, que é o que a documentação de requisições recomenda.

O diamante: duas emissões, uma delas errada

A documentação do combineLatest descreve o comportamento numa frase: sempre que qualquer Observable de entrada emite, ele calcula a saída com os valores mais recentes de todos. Quando as entradas derivam da mesma fonte, uma mudança na fonte vira duas emissões, e a primeira combina um valor novo com um velho.

No exemplo, subtotal$ e total$ derivam do mesmo carrinho. Num estado consistente, total - subtotal é sempre o frete, 10. O script conta quantas emissões chegam e quantas violam essa regra, e repete o mesmo diamante com signal e computed do Angular, que rodam no Node sem aplicação nenhuma:

import { BehaviorSubject, combineLatest, map } from 'rxjs';
import { computed, signal } from '@angular/core';

const N = 1000;

// RxJS: dois valores derivados da mesma fonte, combinados de novo.
const cart$ = new BehaviorSubject({ subtotal: 100, shipping: 10 });
const subtotal$ = cart$.pipe(map((c) => c.subtotal));
const total$ = cart$.pipe(map((c) => c.subtotal + c.shipping));
let emissions = 0, inconsistent = 0;
combineLatest([subtotal$, total$]).subscribe(([subtotal, total]) => {
  emissions++;
  if (total - subtotal !== 10) inconsistent++; // estado consistente: a diferença é o frete
});
emissions = inconsistent = 0; // ignora a emissão inicial
for (let i = 1; i <= N; i++) cart$.next({ subtotal: 100 + i, shipping: 10 });
console.log(`combineLatest: ${N} updates -> ${emissions} emissions, ${inconsistent} inconsistent`);

// Signals do Angular: o mesmo diamante com computed.
const cart = signal({ subtotal: 100, shipping: 10 });
const subtotal = computed(() => cart().subtotal);
const total = computed(() => cart().subtotal + cart().shipping);
let evaluations = 0;
inconsistent = 0;
const summary = computed(() => {
  evaluations++;
  if (total() - subtotal() !== 10) inconsistent++;
  return `${subtotal()} / ${total()}`;
});
summary();
evaluations = 0;
for (let i = 1; i <= N; i++) { cart.set({ subtotal: 100 + i, shipping: 10 }); summary(); }
console.log(`computed, read after each set: ${N} updates -> ${evaluations} evaluations, ${inconsistent} inconsistent`);
evaluations = 0;
for (let i = 1; i <= N; i++) cart.set({ subtotal: 5000 + i, shipping: 10 });
summary();
console.log(`computed, read once at the end: ${N} updates -> ${evaluations} evaluation`);
combineLatest: 1000 updates -> 2000 emissions, 1000 inconsistent
computed, read after each set: 1000 updates -> 1000 evaluations, 0 inconsistent
computed, read once at the end: 1000 updates -> 1 evaluation

Mil atualizações do carrinho viraram 2.000 emissões, e metade delas trazia um total que não batia com o subtotal. O efeito na tela eu não medi: as duas emissões acontecem na mesma tarefa síncrona, antes de o Angular renderizar. O que dá para afirmar é o que vem depois do combineLatest: tudo ali roda duas vezes por mudança (um tap que grava, um log, um switchMap que monta uma consulta), e uma dessas execuções recebe valores que nunca existiram juntos.

A correção em RxJS é não criar o diamante: derivar os dois valores no mesmo map, a partir da mesma emissão. Rodei a variante cart$.pipe(map((c) => [c.subtotal, c.subtotal + c.shipping])) com as mesmas mil atualizações: mil emissões, nenhuma inconsistente.

Com signals, o diamante não produz leitura inconsistente: mil avaliações quando o computed é lido depois de cada set, nenhuma errada, e uma avaliação só quando ele é lido uma vez depois de mil sets. Isso bate com o que o guia de signals descreve: o valor de um computed fica em cache e só é recalculado na próxima leitura depois que uma dependência mudou. No post de julho eu afirmei que computed é glitch-free; aqui está a medição que faltava ali.

distinctUntilChanged ajuda, com uma condição

O outro jeito de gerar trabalho a mais é emitir um valor igual ao anterior. O distinctUntilChanged existe para isso, e a documentação diz como ele compara: com o comparador passado ou, sem comparador, com ===. Essa última parte decide se ele serve para alguma coisa:

import { BehaviorSubject, distinctUntilChanged, map } from 'rxjs';

const N = 1000;
function count(label, buildPipe, updates) {
  const source$ = new BehaviorSubject(updates(0));
  let runs = 0;
  buildPipe(source$).subscribe(() => runs++);
  runs = 0; // ignora a emissão inicial
  for (let i = 1; i <= N; i++) source$.next(updates(i));
  console.log(`${label.padEnd(52)} ${N} updates -> ${runs} downstream runs`);
}

const price = (i) => 101 + (i % 100); // sempre "caro"
count('boolean, no distinctUntilChanged', (p$) => p$.pipe(map((p) => p > 100)), price);
count('boolean + distinctUntilChanged()', (p$) => p$.pipe(map((p) => p > 100), distinctUntilChanged()), price);

const state = (i) => ({ user: { name: 'Ada' }, lastSeen: i }); // o nome nunca muda
count('view model + distinctUntilChanged()',
  (s$) => s$.pipe(map((s) => ({ name: s.user.name })), distinctUntilChanged()), state);
count('view model + distinctUntilChanged((a, b) => a.name === b.name)',
  (s$) => s$.pipe(map((s) => ({ name: s.user.name })), distinctUntilChanged((a, b) => a.name === b.name)), state);
boolean, no distinctUntilChanged                     1000 updates -> 1000 downstream runs
boolean + distinctUntilChanged()                     1000 updates -> 0 downstream runs
view model + distinctUntilChanged()                  1000 updates -> 1000 downstream runs
view model + distinctUntilChanged((a, b) => a.name === b.name) 1000 updates -> 0 downstream runs

Com um booleano, o operador corta as mil execuções para zero. Com um objeto montado num map, o distinctUntilChanged() sem argumento não corta nada: cada map cria um objeto novo, e === entre objetos diferentes é sempre falso. Ele vira um operador que custa e não faz nada. Com o comparador, volta a zero.

O mesmo cuidado vale do lado dos signals. A igualdade padrão de um signal é referencial (Object.is, segundo o guia), então um objeto novo com o mesmo conteúdo conta como mudança.

Typeahead: switchMap, mergeMap, concatMap e exhaustMap

O último caso é concorrência. Uma busca recebe sete teclas, uma a cada 70 ms, até formar "angular". Cada consulta tem uma latência fixa no servidor simulado, escolhida para gerar corrida: a segunda ("an") é a mais lenta, 460 ms, e a última ("angular") é a mais rápida, 100 ms. O script usa tempo virtual (TestScheduler.run), então imprime exatamente os mesmos números em toda execução.

// Mesmas teclas e mesmas latências passando por cada operador de achatamento.
import { NEVER, Observable, concatMap, debounceTime, exhaustMap, map, merge, mergeMap, switchMap, timer } from 'rxjs';
import { TestScheduler } from 'rxjs/testing';

const QUERIES = ['a', 'an', 'ang', 'angu', 'angul', 'angula', 'angular'];
const TYPED_AT = [0, 70, 140, 210, 280, 350, 420]; // uma tecla a cada 70 ms
const LATENCY = { a: 300, an: 460, ang: 200, angu: 350, angul: 150, angula: 250, angular: 100 };

function run(label, flatten) {
  const scheduler = new TestScheduler(() => {});
  const stats = { started: 0, cancelled: 0, delivered: 0, inFlight: 0, maxInFlight: 0, shown: [] };

  scheduler.run(() => {
    const search = (q) =>
      new Observable((subscriber) => {
        stats.started++;
        stats.maxInFlight = Math.max(stats.maxInFlight, ++stats.inFlight);
        let done = false;
        const response = timer(LATENCY[q]).subscribe(() => {
          done = true;
          stats.inFlight--;
          subscriber.next(q);
          subscriber.complete();
        });
        return () => {
          if (!done) { stats.cancelled++; stats.inFlight--; }
          response.unsubscribe();
        };
      });

    // NEVER mantém a entrada aberta, como um campo de texto de verdade (o debounceTime emite antes se a fonte completa).
    const keystrokes$ = merge(...QUERIES.map((q, i) => timer(TYPED_AT[i]).pipe(map(() => q))), NEVER);
    flatten(keystrokes$, search).subscribe((q) => {
      stats.delivered++;
      stats.shown.push(`${q}@${scheduler.now()}ms`);
    });
  });

  const last = stats.shown.at(-1);
  console.log(
    `${label.padEnd(28)} sent ${stats.started} | not sent ${QUERIES.length - stats.started} | ` +
    `cancelled ${stats.cancelled} | max in flight ${stats.maxInFlight} | ` +
    `screen ends on "${last}"${last.startsWith('angular@') ? '' : '  <-- stale'}`
  );
}

run('mergeMap', (k$, s) => k$.pipe(mergeMap(s)));
run('concatMap', (k$, s) => k$.pipe(concatMap(s)));
run('exhaustMap', (k$, s) => k$.pipe(exhaustMap(s)));
run('switchMap', (k$, s) => k$.pipe(switchMap(s)));
run('debounceTime(200)+switchMap', (k$, s) => k$.pipe(debounceTime(200), switchMap(s)));
mergeMap                     sent 7 | not sent 0 | cancelled 0 | max in flight 5 | screen ends on "angula@600ms"  <-- stale
concatMap                    sent 7 | not sent 0 | cancelled 0 | max in flight 1 | screen ends on "angular@1810ms"
exhaustMap                   sent 2 | not sent 5 | cancelled 0 | max in flight 1 | screen ends on "angula@600ms"  <-- stale
switchMap                    sent 7 | not sent 0 | cancelled 6 | max in flight 1 | screen ends on "angular@520ms"
debounceTime(200)+switchMap  sent 1 | not sent 6 | cancelled 0 | max in flight 1 | screen ends on "angular@720ms"

mergeMap chegou a ter cinco requests em voo ao mesmo tempo, e as respostas foram exibidas na ordem em que o servidor terminou. A tela termina em "angula", aos 600 ms: a resposta certa chegou aos 520 ms e foi sobrescrita por uma mais lenta. Com um backend local respondendo tudo em tempos parecidos, essa corrida pode nunca aparecer em desenvolvimento.

exhaustMap ignora valores novos enquanto o anterior não terminou (é a definição na documentação). Mandou dois requests e também terminou em "angula". Errado para busca, certo para um botão de enviar, onde o segundo clique deve ser ignorado.

concatMap acerta, mas leva 1.810 ms: a soma de todas as latências, porque cada request espera o anterior. A documentação ainda avisa que, se os valores chegam mais rápido do que os internos completam, eles se acumulam num buffer sem limite.

switchMap acerta em 520 ms. Mas ele enviou sete requests e cancelou seis. O cancelamento acontece no cliente: o HttpClient despacha a requisição no subscribe e aborta no unsubscribe, como diz a documentação, mas abortar não desfaz o que o servidor já recebeu.

debounceTime(200) + switchMap mandou um request só e acertou, em 720 ms. A troca é explícita: 200 ms a mais de latência por seis requests a menos. Qual dos dois vale mais depende do custo da consulta no backend.

Uma armadilha em que eu caí escrevendo o teste: na primeira versão, o fluxo de teclas completava depois da última tecla, e o debounceTime emitia na hora, sem esperar os 200 ms. A documentação do debounceTime descreve isso: se a fonte completa durante a espera, o último valor guardado é emitido antes da conclusão. Um campo de texto de verdade nunca completa; o NEVER no merge reproduz isso.

No Angular 22, o rxResource é estável desde a 22.0 e entrega a semântica do switchMap sem escrevê-lo. O guia de resources diz que um resource aborta o carregamento em andamento quando os parâmetros mudam, e no código do 22.2.1 esse aborto faz unsubscribe do Observable anterior. Se os parâmetros devolvem undefined, o carregamento nem começa e o status fica idle:

import { Component, inject, signal } from '@angular/core';
import { HttpClient } from '@angular/common/http';
import { rxResource } from '@angular/core/rxjs-interop';

interface Product {
  id: string;
  name: string;
}

@Component({ selector: 'app-product-search', template: `...` })
export class ProductSearch {
  private http = inject(HttpClient);
  query = signal('');

  results = rxResource({
    // string vazia vira undefined: sem request, status 'idle'
    params: () => this.query() || undefined,
    stream: ({ params }) => this.http.get<Product[]>('/api/products', { params: { q: params } }),
  });
}

O debounce continua por sua conta. E um detalhe que a tipagem do 22.2.1 documenta: se o Observable completar sem emitir nada (um catchError(() => EMPTY), por exemplo), o Angular lança o erro NG0991, porque o resource fica sem valor e sem erro para mostrar. Quem migrar um typeahead que engolia erros com catchError(() => EMPTY) dentro do switchMap precisa devolver um valor (of([])) ou deixar o erro chegar ao resource.

Por que isso importa

As medições formam uma escada. Um operador custa nanossegundos. Uma assinatura custa meio microssegundo e pouco mais de 3 KB. O que custa milissegundos, megabytes ou um número errado na tela é sempre trabalho: repetido (três requests no lugar de um), que não para (mil timers depois que todo mundo saiu), emitido a mais (duas emissões do diamante para cada mudança) ou na ordem errada (o mergeMap mostrando "angula").

Nenhum desses problemas aparece num profiler como "RxJS lento". Eles aparecem como requests duplicados na aba de rede, memória que sobe a cada navegação, um total que pisca e uma busca que mostra resultado velho. Por isso, quando investigo performance com RxJS, eu conto em vez de cronometrar: quantas vezes o produtor rodou, quantos teardowns aconteceram, quantas emissões chegaram. Contar é barato, dá o mesmo resultado em toda execução e vale igual no Node e no navegador. Nanossegundos não valem.

Referências

Comentários

Carregando comentários...