O Amazon Kinesis no Node.js permite publicar e processar streams de eventos em tempo real usando um serviço gerenciado da AWS. Kinesis Data Streams organiza registros em shards, preserva ordem por partition key e mantém dados por um período configurável para múltiplos consumidores independentes.
O serviço é útil para clickstream, telemetria, eventos de domínio, logs, métricas e pipelines de dados. Porém, ele não garante processamento exatamente uma vez. Producers podem reenviar registros e consumers podem receber novamente após falhas, portanto idempotência e checkpoints fazem parte do design.
Neste guia, você aprenderá a usar o AWS SDK v3, escolher partition keys, enviar lotes, consumir com Lambda ou polling, tratar partial batch failure, monitorar iterator age, escalar shards, aplicar segurança e fazer graceful shutdown.
O que é Kinesis Data Streams?
A documentação oficial de Amazon Kinesis Data Streams descreve um serviço para coletar e processar grandes fluxos de registros em tempo real. Várias aplicações podem consumir o mesmo stream de forma independente para analytics, alertas, armazenamento e processamento operacional.
O atraso entre publicar e disponibilizar um registro costuma ser baixo, mas depende de capacidade, consumers, rede e configuração.
Conceitos principais
- Stream: conjunto lógico de registros.
- Shard: unidade de capacidade e ordering.
- Partition key: define o shard de destino.
- Sequence number: posição do registro no shard.
- Producer: publica registros.
- Consumer: lê e processa.
- Retention: período em que registros permanecem disponíveis.
Instalando o AWS SDK v3
npm install @aws-sdk/client-kinesisO SDK v3 é modular. Importe apenas os clients e commands necessários.
Criando o client
import { KinesisClient } from '@aws-sdk/client-kinesis';
export const kinesis = new KinesisClient({
region: process.env.AWS_REGION,
requestHandler: {
requestTimeout: 5000
}
});Em EC2, ECS, EKS ou Lambda, prefira credenciais temporárias por IAM Role. Não coloque access keys no código ou imagem.
Enviando um registro
import { PutRecordCommand } from '@aws-sdk/client-kinesis';
const event = {
id: crypto.randomUUID(),
type: 'order.created',
source: 'orders-api',
version: 1,
occurredAt: new Date().toISOString(),
data: {
orderId,
customerId,
totalCents
}
};
const response = await kinesis.send(new PutRecordCommand({
StreamName: process.env.KINESIS_STREAM,
PartitionKey: orderId,
Data: Buffer.from(JSON.stringify(event))
}));Use uma chave que preserve a ordem necessária e distribua carga. orderId mantém eventos do mesmo pedido juntos sem concentrar todos em um único shard.
Escolhendo partition key
Boas opções dependem do domínio:
orderIdpara eventos de pedido;deviceIdpara telemetria;customerIdpara atividade do cliente;sessionIdpara clickstream.
Uma chave com poucos valores cria hotspots. Uma chave aleatória distribui bem, mas perde ordering por entidade.
Ordering
Kinesis preserva ordem apenas dentro de um shard. Registros com a mesma partition key tendem ao mesmo shard, mas consumers precisam processar sequencialmente quando a ordem importa.
Não existe uma ordem global simples entre todos os shards.
Explicit hash key
O producer pode fornecer ExplicitHashKey, mas isso normalmente é desnecessário e aumenta risco de distribuição ruim. Prefira partition key bem escolhida.
Envio em lote
A documentação de exemplos Kinesis com AWS SDK v3 mostra o uso de PutRecords:
import { PutRecordsCommand } from '@aws-sdk/client-kinesis';
const records = events.map(event => ({
PartitionKey: event.data.orderId,
Data: Buffer.from(JSON.stringify(event))
}));
const result = await kinesis.send(new PutRecordsCommand({
StreamName: process.env.KINESIS_STREAM,
Records: records
}));Um request bem-sucedido pode conter falhas individuais. Verifique FailedRecordCount e cada item.
Retry apenas dos registros falhos
const failed = result.Records
.map((record, index) => ({ record, input: records[index] }))
.filter(item => item.record.ErrorCode);
for (const item of failed) {
await retryWithJitter(item.input);
}Não repita todo o lote, pois registros já aceitos seriam duplicados.
Backoff e jitter
Erros de throughput e indisponibilidade temporária devem usar backoff exponencial com jitter. Limite tentativas e tempo total.
const delay = Math.min(100 * 2 ** attempt, 5000);
const jitter = Math.random() * delay;
await new Promise(resolve => setTimeout(resolve, jitter));Idempotência do producer
Uma conexão pode falhar depois de o serviço aceitar o registro. O producer não sabe se deve repetir. Inclua event.id estável e faça deduplicação no consumer.
Consulte Idempotência em APIs Node.js.
Outbox Pattern
Publicar após uma transação PostgreSQL pode falhar e perder o evento. Use Outbox:
- grave a entidade e o evento na mesma transação;
- um relay lê eventos pendentes;
- publica no Kinesis com ID estável;
- marca o evento como enviado;
- retries reutilizam o mesmo ID.
Veja Outbox Pattern no Node.js.
CloudEvents
Um envelope CloudEvents ajuda a padronizar id, source, type, time e subject. Consulte CloudEvents no Node.js.
Limites de payload
Registros possuem limite de tamanho. Não envie arquivos, imagens ou objetos enormes. Armazene o conteúdo em S3 e envie um identificador, hash e metadata.
On-demand e provisioned
No modo on-demand, a AWS ajusta capacidade automaticamente dentro das regras do serviço. No modo provisioned, a equipe gerencia shards.
On-demand reduz operação inicial. Provisioned pode oferecer controle de custo e capacidade em workloads previsíveis.
Shards e throughput
Cada shard possui limites de leitura e escrita. Monitore:
- bytes por segundo;
- records por segundo;
- throttling;
- hot partition keys;
- iterator age;
- consumers por shard.
Mais shards não corrigem uma partition key concentrada.
Enhanced fan-out
Enhanced fan-out fornece throughput de leitura dedicado por consumer registrado e usa HTTP/2 push. É útil quando vários consumidores precisam de baixa latência sem compartilhar o limite de leitura do shard.
Ele possui custo e exige registro de consumer.
Consumo com Lambda
Lambda pode receber lotes do stream:
export async function handler(event) {
for (const record of event.Records) {
const payload = JSON.parse(
Buffer.from(record.kinesis.data, 'base64').toString('utf8')
);
await processEvent(payload);
}
}A integração gerencia polling e checkpoints. Uma falha pode fazer o lote ou parte dele ser processado novamente.
Partial batch response
A AWS documenta retorno de sequence numbers falhos:
export async function handler(event) {
const batchItemFailures = [];
for (const record of event.Records) {
try {
const payload = JSON.parse(
Buffer.from(record.kinesis.data, 'base64').toString('utf8')
);
await processEvent(payload);
} catch (error) {
batchItemFailures.push({
itemIdentifier: record.kinesis.sequenceNumber
});
break;
}
}
return { batchItemFailures };
}Como streams preservam ordem no shard, o consumer deve considerar o primeiro registro falho e os posteriores.
Bisect batch on function error
A configuração do event source mapping pode dividir lotes falhos para localizar poison messages. Combine com limite de retries e destino de falha.
Parallelization factor
Lambda pode processar múltiplos lotes por shard, mantendo ordem por partition key. Aumentar paralelismo eleva throughput, mas também conexões e concorrência em dependências.
Consumer próprio
Um consumer de longa duração pode usar GetShardIterator e GetRecords. Porém, precisa descobrir shards, renovar iterators, controlar checkpoints, rebalancear workers e tratar resharding.
Para aplicações complexas, a Kinesis Client Library gerencia leases e checkpoints. O suporte oficial principal é em outros runtimes; em Node.js, avalie Lambda, serviços gerenciados ou implementação cuidadosamente testada.
Checkpoint
Checkpoint registra até qual sequence number o consumer concluiu. Faça checkpoint somente depois de persistir efeitos. Se registrar antes, uma falha perde o processamento.
Inbox Pattern
Para efeitos PostgreSQL:
BEGIN;
INSERT INTO processed_events(consumer_name, event_id)
VALUES ($1, $2)
ON CONFLICT DO NOTHING;
UPDATE projections
SET total = total + $3
WHERE id = $4;
COMMIT;Se o evento já existe, não repita a atualização. Veja Inbox Pattern no Node.js.
Poison messages
Uma mensagem inválida pode bloquear a progressão do shard. Defina:
- validação de schema;
- limite de retries;
- destino de falha;
- alerta;
- processo de correção e replay.
Schema evolution
Inclua version e mantenha consumers compatíveis com eventos antigos. Adicione campos opcionais, preserve significado e teste replay.
Retention
Aumentar retenção permite replay mais longo, mas eleva custo. Escolha conforme tempo de recuperação e requisitos de auditoria.
Replay
Um novo consumer pode começar no início da retenção ou em um timestamp. Efeitos precisam ser idempotentes, pois replay deliberadamente processa registros antigos.
Resharding
Split e merge alteram shards. Consumers precisam descobrir os novos shards e respeitar relações parent-child para não perder ordenação.
Segurança IAM
Um producer precisa apenas de ações como:
{
"Effect": "Allow",
"Action": [
"kinesis:PutRecord",
"kinesis:PutRecords"
],
"Resource": "arn:aws:kinesis:sa-east-1:123456789012:stream/order-events"
}Consumers recebem apenas leitura e descrição necessárias. Evite kinesis:*.
Criptografia
Kinesis oferece criptografia em trânsito e server-side encryption. Use KMS quando a política exigir controle de chaves e configure IAM para uso da key.
VPC endpoints
Um interface endpoint pode manter tráfego dentro da rede AWS e reduzir dependência de internet pública. Configure security groups, DNS e policies.
Dados sensíveis
Não envie tokens, senhas ou dados pessoais sem necessidade. A retenção e múltiplos consumers ampliam exposição. Aplique minimização, criptografia e controle de acesso.
Graceful shutdown
Producers devem concluir lotes pendentes e fechar o client. Consumers próprios precisam parar polling, concluir registros e salvar checkpoint.
Veja Graceful Shutdown no Node.js.
Observabilidade
Monitore:
- PutRecord e PutRecords latency;
- failed records;
- write e read throttling;
- IncomingBytes e IncomingRecords;
- GetRecords.IteratorAgeMilliseconds;
- Lambda errors e retries;
- partial batch failures;
- hot partition keys;
- idade do evento ao processar.
Alertas
Iterator age crescente indica consumer atrasado. Combine com erro, concorrência, throttling e capacidade. Um alerta apenas por volume pode disparar durante um pico saudável.
Testes
Teste:
- lote com falhas parciais;
- producer timeout após aceite;
- consumer crash antes do checkpoint;
- poison message;
- hot partition key;
- resharding;
- replay;
- perda temporária de credenciais.
Kinesis ou Kafka?
Kinesis é gerenciado e integrado à AWS. Kafka possui ecossistema e modelo operacional diferentes. A escolha depende de portabilidade, features, equipe, throughput e custo.
Consulte Kafka com Node.js.
Kinesis ou SQS?
SQS é uma fila em que mensagens são distribuídas entre consumidores. Kinesis mantém um log ordenado por shard e permite múltiplos consumers independentes e replay.
Veja Amazon SQS no Node.js.
Erros comuns
- Partition key constante: um shard recebe toda a carga.
- Repetir lote inteiro: registros aceitos são duplicados.
- Checkpoint antes do efeito: falha perde processamento.
- Consumer sem idempotência: retry duplica atualizações.
- Payload grande: custo e throughput degradam.
- Sem limite de poison message: shard fica bloqueado.
- IAM amplo: serviço acessa streams indevidos.
- Ignorar iterator age: backlog cresce silenciosamente.
Conclusão
O Amazon Kinesis no Node.js oferece streams gerenciados com shards, ordering por partition key, retenção e múltiplos consumers. O AWS SDK v3 permite publicar registros individuais ou em lote, enquanto Lambda simplifica consumo e checkpoints.
Distribua partition keys, trate falhas parciais e use IDs estáveis. Aplique Inbox e Outbox, limite retries e monitore iterator age e throttling. Com IAM mínimo e deploy idempotente, Kinesis sustenta pipelines em tempo real sem esconder as garantias de entrega pelo menos uma vez.



