O Apache Pulsar no Node.js oferece mensageria e streaming distribuído com topics particionados, múltiplos modelos de subscription, retenção, backlog e separação entre compute e storage. A biblioteca oficial permite criar producers, consumers e readers em aplicações JavaScript e TypeScript.
Pulsar é adequado para eventos, filas, pipelines e comunicação entre serviços. Ele não elimina idempotência, schema evolution ou observabilidade. Consumers podem receber mensagens novamente após falhas, e um cluster mal dimensionado pode acumular backlog ou storage rapidamente.
Neste guia, você aprenderá a instalar o client, produzir e consumir mensagens, escolher subscription type, aplicar ACK, retry, DLQ, batching, partitioning, schemas, TLS e graceful shutdown.
O que é Apache Pulsar?
A documentação oficial do Pulsar Node.js client informa que a biblioteca cria producers, consumers e readers. Pulsar é uma plataforma distribuída de messaging e streaming da Apache Software Foundation.
A arquitetura separa brokers do armazenamento persistente, normalmente Apache BookKeeper, permitindo escalar serving e storage independentemente.
Instalação
npm install pulsar-clientO pacote inclui bindings nativos. Confirme compatibilidade com Node.js, sistema operacional e arquitetura antes do build da imagem.
Criando o client
import Pulsar from 'pulsar-client';
const client = new Pulsar.Client({
serviceUrl: process.env.PULSAR_SERVICE_URL,
operationTimeoutSeconds: 30,
ioThreads: 4,
messageListenerThreads: 4
});Use um client compartilhado por processo, em vez de criar um para cada mensagem.
Producer
const producer = await client.createProducer({
topic: 'persistent://commerce/orders/order-events',
producerName: 'orders-api',
sendTimeoutMs: 5000,
batchingEnabled: true,
batchingMaxMessages: 500,
compressionType: 'LZ4'
});O topic inclui tenant, namespace e nome. Namespaces ajudam a aplicar retenção, quotas, autenticação e políticas.
Enviando mensagem
const messageId = await producer.send({
data: Buffer.from(JSON.stringify(event)),
properties: {
eventType: 'order.created',
schemaVersion: '1',
traceId
},
partitionKey: event.orderId,
eventTimestamp: Date.now()
});partitionKey mantém mensagens da mesma chave na mesma partição, ajudando ordering. Não use uma chave global, pois cria hotspot.
Payloads pequenos
Envie evento com IDs e dados necessários, não snapshots gigantes. Mensagens grandes aumentam rede, memória, retenção e tempo de retry.
Producer batching
Batching combina várias mensagens em uma requisição. Melhora throughput, mas adiciona pequena latência. Ajuste tamanho e atraso conforme o SLO.
Compression
Compressão reduz rede e storage, mas usa CPU. Meça LZ4, ZSTD ou algoritmo suportado com payloads reais.
Consumer básico
const consumer = await client.subscribe({
topic: 'persistent://commerce/orders/order-events',
subscription: 'billing-service',
subscriptionType: 'Shared',
ackTimeoutMs: 60_000,
receiverQueueSize: 1000
});
while (!shuttingDown) {
const message = await consumer.receive();
try {
const event = eventSchema.parse(
JSON.parse(message.getData().toString())
);
await handleEvent(event);
consumer.acknowledge(message);
} catch (error) {
consumer.negativeAcknowledge(message);
}
}Subscription Exclusive
Somente um consumer pode conectar. É útil quando uma única instância precisa preservar processamento serial. Falha de conexão permite outro assumir.
Subscription Shared
Mensagens são distribuídas entre vários consumers. A ordem global não é garantida. Use para jobs independentes e alto throughput.
Failover
Vários consumers conectam, mas um fica ativo por partição e outros aguardam. Ao falhar, o próximo assume. É útil para ordered consumption com standby.
Key_Shared
Mensagens com a mesma key são entregues ao mesmo consumer, enquanto chaves diferentes são paralelizadas. Producers precisam desativar batching incompatível ou usar key-based batching conforme a versão.
ACK individual
Confirma uma mensagem específica:
consumer.acknowledge(message);Use quando processamento é independente.
Cumulative ACK
Confirma a mensagem e todas as anteriores. Não é compatível com todos os tipos de subscription. Use somente quando a ordem e o processamento sequencial garantem segurança.
Negative acknowledgement
consumer.negativeAcknowledge(message);A mensagem volta após o atraso configurado. Falhas permanentes podem causar loop infinito; use retry letter topic e DLQ.
Ack timeout
Se a mensagem não é confirmada dentro do prazo, pode ser redeliver. O timeout deve superar o p99 do processamento com margem. Jobs longos podem usar outra estratégia.
Idempotência
Um consumer pode executar o efeito e falhar antes do ACK. Ao receber novamente, precisa reconhecer que já concluiu:
INSERT INTO processed_events(consumer, event_id, processed_at)
VALUES ($1, $2, now())
ON CONFLICT DO NOTHING
RETURNING event_id;Se nenhuma linha for retornada, confirme a mensagem sem repetir o efeito. Veja Inbox Pattern no Node.js.
Dead Letter Policy
const consumer = await client.subscribe({
topic,
subscription: 'billing-service',
subscriptionType: 'Shared',
deadLetterPolicy: {
maxRedeliverCount: 5,
deadLetterTopic: 'persistent://commerce/dlq/billing-events'
}
});Monitore a DLQ, registre motivo e crie processo de reprocessamento. Não deixe mensagens acumularem silenciosamente.
Retry letter topic
Pulsar suporta retry topics em combinações específicas. Eles separam mensagens em espera do backlog principal e permitem backoff. Confirme suporte no client e na versão do broker.
Delayed delivery
await producer.send({
data: payload,
deliverAfter: 60_000
});A mensagem fica indisponível até o prazo. Delayed delivery não substitui workflow durável complexo.
Topics particionados
Crie múltiplas partições para aumentar throughput:
pulsar-admin topics create-partitioned-topic \
persistent://commerce/orders/order-events \
--partitions 12A quantidade deve considerar producers, consumers, ordering e capacidade. Mais partições aumentam metadata e conexões.
Ordering
Ordering é por partição ou key, não global. Se eventos do mesmo pedido precisam de ordem, use partitionKey=orderId e Key_Shared ou consumer compatível.
Deduplicação do broker
Pulsar pode deduplicar mensagens de producer usando sequência e nome do producer. Isso reduz duplicatas de envio, mas não substitui idempotência do consumer.
Producer send timeout
Quando o timeout ocorre, a mensagem pode ter sido persistida. Repetir com um novo ID pode duplicar. Use event ID estável e deduplicação.
Schema
Pulsar oferece schema registry integrado para JSON, Avro, Protobuf e outros formatos conforme client. Mesmo com JSON manual, mantenha eventType e schemaVersion.
Veja CloudEvents no Node.js para envelope padronizado.
Schema evolution
Adicione campos opcionais, preserve semântica e teste consumidores antigos. Remover ou alterar tipo pode quebrar replay de mensagens históricas.
Reader
const reader = await client.createReader({
topic,
startMessageId: Pulsar.MessageId.earliest()
});
while (await reader.hasNext()) {
const message = await reader.readNext();
await inspect(message);
}Reader não usa subscription cursor da mesma forma que consumer. É útil para replay e auditoria, mas precisa controlar posição.
Retention e backlog
Retention preserva mensagens mesmo após ACK, conforme política. Backlog mantém mensagens não confirmadas. Configure tamanho e tempo para evitar storage ilimitado.
Backlog quota
Quando o backlog supera a quota, o namespace pode bloquear producers ou descartar dados conforme configuração. Escolha a política conscientemente e alerte antes do limite.
TTL
Message TTL expira mensagens não confirmadas. Não use para tarefas que não podem ser perdidas. Para notificações obsoletas, pode ser apropriado.
Multi-tenancy
Tenants e namespaces isolam políticas e autenticação. Não use tenant enviado pelo cliente para construir topic sem allowlist, pois pode permitir acesso cruzado.
Consulte Multi-Tenancy no Node.js.
Autenticação por token
const client = new Pulsar.Client({
serviceUrl,
authentication: new Pulsar.AuthenticationToken({
token: process.env.PULSAR_TOKEN
})
});Use token de curta duração ou mecanismo gerenciado quando disponível.
TLS
const client = new Pulsar.Client({
serviceUrl: 'pulsar+ssl://pulsar.internal:6651',
tlsTrustCertsFilePath: '/etc/pulsar/ca.pem',
tlsValidateHostname: true,
tlsAllowInsecureConnection: false
});Valide hostname e CA. Não desative validação em produção.
Outbox Pattern
Para publicar depois de uma transação PostgreSQL:
- grave entidade e outbox no mesmo COMMIT;
- relay publica no Pulsar;
- usa event ID e producer name estáveis;
- marca o registro como enviado;
- consumer aplica Inbox.
Veja Outbox Pattern no Node.js.
Transactions Pulsar
Pulsar oferece transactions em versões e clients compatíveis para produzir e confirmar mensagens atomicamente dentro da plataforma. Isso não inclui automaticamente seu banco PostgreSQL.
Graceful shutdown
async function shutdown() {
shuttingDown = true;
await consumer.close();
await producer.flush();
await producer.close();
await client.close();
}
process.on('SIGTERM', shutdown);Interrompa novas mensagens, conclua as ativas e só então feche. Consulte Graceful Shutdown no Node.js.
Observabilidade
Monitore:
- publish e consume rate;
- backlog e idade da mensagem mais antiga;
- redelivery;
- unacked messages;
- DLQ;
- latência de publish e ACK;
- storage e bookies;
- conexões e erros do client;
- partições desequilibradas.
Testes
Use um cluster real em container ou ambiente isolado. Teste producer timeout, consumer crash antes do ACK, redelivery, DLQ, ordering por key, schema incompatível e shutdown.
Erros comuns
- Consumer sem idempotência: redelivery duplica efeitos.
- Shared esperando ordem: mensagens são paralelizadas.
- Partition key global: uma partição vira hotspot.
- Ack antes do COMMIT: falha perde o efeito.
- Nack infinito: poison message consome recursos.
- Backlog sem quota: storage cresce sem limite.
- TLS inseguro: broker não é autenticado.
- Payload grande: throughput e retenção degradam.
Conclusão
O Apache Pulsar no Node.js oferece producers, consumers, subscriptions, partições, retenção e streaming distribuído. O client oficial permite escolher Exclusive, Shared, Failover e Key_Shared conforme ordering e escala.
Use event IDs, Inbox e Outbox, limite retries e monitore backlog e DLQ. Configure TLS, namespaces e quotas e teste redelivery. Assim, Pulsar sustenta pipelines de eventos sem transformar entrega pelo menos uma vez em efeitos duplicados.


