Lab 29 — Streaming de mudança do banco
O problema, e a empresa que o tem
O L16 deixou a Cadência com uma tabela única, pedidos_cadencia, servindo quatro padrões de acesso sem nenhum Scan. O problema agora é outro: três sistemas diferentes PRECISAM saber quando um pedido muda de status, e hoje nenhum deles sabe até perguntar.
O mais visível é o painel de operação que os gerentes de loja olham no início do turno. Ele é alimentado por uma função agendada a cada 2 minutos, que compara o campo `updated_at` com o horário da última execução — o jeito mais natural de fazer isso para quem vem de banco relacional, e funcionou bem em teste.
Em produção, apareceu um problema que ninguém tinha previsto: quando um pedido é separado e, minutos depois, cancelado — cliente desistiu, fraude, endereço inválido — o time de estoque precisa ser avisado para devolver o item já embalado à prateleira. Essa regra depende de saber que o pedido PASSOU por "separado" antes de chegar a "cancelado". A função de polling só vê o status final: se as duas mudanças aconteceram dentro do mesmo intervalo de 2 minutos, o alerta simplesmente nunca dispara — não atrasado, nunca. Uma auditoria de estoque no mês anterior encontrou 30 casos assim, todos descobertos só na contagem física.
O que este laboratório NÃO é
Não é uma fila de mensagens de uso geral: DynamoDB Streams está atrelado à tabela que o gerou, retém só 24 horas e serve bem UM padrão — reagir à mudança de um item específico. Quando o requisito vira "vários sistemas diferentes, cada um com seu próprio ritmo de consumo, precisando reprocessar semanas de histórico", a resposta certa é Kinesis Data Streams na frente do DynamoDB, e isso é o L63. Este módulo também não trata CDC de banco relacional — WAL, slot de replicação, DMS — porque isso já foi construído no L19, sobre Aurora/RDS; aqui a origem é DynamoDB, e o mecanismo nativo é outro.
O que você vai conseguir fazer
Objetivos verificáveis: cada um se prova com uma medição na seção de implantação, não com a sensação de ter entendido a teoria.
- Explicar por que uma consulta periódica por timestamp PERDE transição, e não apenas atrasa.
- Ligar DynamoDB Streams numa tabela existente com o tipo de visualização correto para o caso de uso.
- Distinguir o que NEW_IMAGE sozinho responde do que só NEW_AND_OLD_IMAGES responde.
- Escrever um consumidor Lambda via event source mapping, com relato parcial de falha de lote.
- Provar, com medição, que a ordem de processamento é preservada DENTRO da mesma chave de partição.
- Explicar por que a ordem NÃO é garantida entre chaves de partição diferentes, e quando isso importa.
- Tornar a escrita do consumidor idempotente diante de entrega pelo menos uma vez (at-least-once).
- Reconhecer o limite de retenção de 24 horas e desenhar o plano de reconciliação para quando ele é excedido.
O que a certificação cobra disto
| Conceito | Certificação | Como aparece aqui | O que dominar |
|---|---|---|---|
| Change data capture (CDC) | SAP-C02, DVA-C02 | DynamoDB Streams como CDC nativo da tabela | a diferença entre "consultar mudança" e "ser avisado de mudança" |
| Tipos de visualização do stream | DVA-C02 | NEW_AND_OLD_IMAGES para decidir com base na transição | KEYS_ONLY, NEW_IMAGE, OLD_IMAGE, NEW_AND_OLD_IMAGES — o que cada um contém |
| Ordem por partição, não global | SAP-C02 | shard dedicado por partição da tabela | ordem garantida DENTRO da chave; não garantida ENTRE chaves diferentes |
| Retenção de 24 horas | SAP-C02, DVA-C02 | plano de reconciliação para downtime maior | stream não é fila durável para sempre |
| At-least-once vs exactly-once | SAP-C02 | idempotência na escrita do consumidor | a entrega é pelo menos uma vez; "exatamente uma vez" é responsabilidade da aplicação |
| Event source mapping do Lambda | DVA-C02 | quem faz o polling dos shards é o serviço, não o código | tamanho de lote, janela de batching, bisecção em erro |
| ReportBatchItemFailures | DVA-C02 | falha parcial de lote sem reprocessar tudo | formato de resposta com itemIdentifier = SequenceNumber |
| DynamoDB Streams vs Kinesis Data Streams for DynamoDB | SAP-C02 | quando um consumidor basta, e quando não basta | retenção (24 h vs até 1 ano) e fan-out para múltiplos consumidores independentes |
Onde isto costuma ser cobrado errado
A pergunta clássica descreve um sistema que "consulta a cada X segundos" para saber o que mudou, e pede o problema com a abordagem quando X diminui. A resposta esperada não é "ainda vai ficar lento" — é que reduzir o intervalo não elimina a perda de transição intermediária, só a torna menos provável. A resposta certa é trocar consulta por notificação: CDC, não polling mais frequente.
Requisitos, e como cada um muda o desenho
Requisito que não aparece numa configuração do stream ou do consumidor é intenção. A coluna da direita é onde cada um deixou marca.
| Requisito | Valor declarado | O que ele decide no desenho |
|---|---|---|
| Latência do painel | segundos, não minutos | obriga trocar consulta periódica por consumo de stream — nenhum ajuste de intervalo de polling chega perto disso |
| Nenhuma transição perdida | inclusive dentro de janelas curtas | obriga NEW_AND_OLD_IMAGES e escrita append-only no painel, uma linha por transição |
| Alerta de reposição correto | dispara sempre que SEPARADO precede CANCELADO | só é decidível com a imagem ANTES e DEPOIS no mesmo registro |
| Sem duplicar alerta em reprocessamento | a entrega é pelo menos uma vez | escrita idempotente com SK derivado do SequenceNumber do stream |
| Falha isolada não trava outros pedidos | requisito operacional da equipe pequena | bisecção de lote em erro, mais DLQ para o que não se resolve sozinho |
| Downtime do consumidor tolerado | até poucas horas, não dias | compatível com a retenção de 24 h do stream; acima disso exige reconciliação manual |
| Sem sharding manual, sem tuning contínuo | equipe de duas pessoas, sem plantão | DynamoDB Streams em vez de Kinesis: sem shard para provisionar ou rebalancear |
| Rastreabilidade da decisão | saber quando e por que cada alerta disparou | log de transições no painel, e não só o estado atual |
Arquitetura mínima: consultar em vez de ser avisado
Este é o desenho que a Cadência tem hoje, e ele é legítimo como ponto de partida: usa só o que a tabela já tinha (`updated_at`) e nenhum serviço novo. O laboratório começa medindo a perda, não o atraso — porque atraso todo mundo já suspeita; perda, ninguém tinha percebido até a auditoria de estoque.
- → dispara a cada 2 minutos
- → Query no GSI2: updated_at > última_execução
- → linha final de cada pedido alterado no intervalo
- → sobrescreve o status conhecido do pedido
- → abre o painel e vê o status mais recente
- Fora da AWS
- Integração de apps
- Compute
- Banco de dados
Este desenho publica, e resolve o caso comum: o gerente vê o status do pedido com até dois minutos de atraso, e ninguém reclama. O defeito não aparece na demonstração — aparece quando um pedido muda de status DUAS VEZES dentro do intervalo, porque a consulta enxerga o banco no instante em que pergunta, nunca o caminho que ele percorreu. Percorra os passos e repare que o alerta de reposição de estoque não tem como disparar, não importa a frequência do relógio.
- O relógio dispara, e o banco decide se havia algo. A função roda a cada 2 minutos independentemente de ter ocorrido mudança. O custo desta consulta é proporcional ao INTERVALO configurado, não à quantidade de eventos reais — rodar mais rápido custa mais RCU, não corrige o defeito central.
- A pergunta sempre olha para trás, e o atraso é a própria pergunta. Uma mudança que ocorre logo depois de uma execução espera o intervalo inteiro até a próxima; uma que ocorre pouco antes é vista quase na hora. Em média, o atraso é metade do intervalo — 1 minuto neste desenho — e no pior caso é o intervalo cheio.
- updated_at sobrescreve; não acumula. Se um pedido vai de PENDENTE a SEPARADO a CANCELADO em 40 segundos, todos dentro do mesmo intervalo de 2 minutos, a Query devolve UMA linha, com o status FINAL. Não existe consulta que recupere o SEPARADO intermediário depois que o item foi sobrescrito — a informação não está atrasada, está perdida.
- O painel reflete o banco agora, não o histórico. A função grava por cima do que já sabia. O painel_operacao é um espelho do estado atual — não guarda nenhuma trilha de transição, porque nunca recebeu essa informação para guardar.
- O alerta que deveria disparar, nunca tem como disparar. A regra de negócio "avisar o estoque quando um pedido é cancelado DEPOIS de separado" precisa saber o ANTES e o DEPOIS da mudança. Este desenho só tem o DEPOIS. Não há frequência de polling que resolva isso: mesmo consultando a cada segundo, uma leitura de snapshot nunca contém a transição, só o estado.
- Por que alguém constrói assim. Porque a tabela já tem `updated_at`, e comparar timestamp é o padrão que quem vem do relacional já sabe fazer — nenhum serviço novo, nenhum conceito de stream ou de shard. Funciona bem enquanto a mudança é rara comparada ao intervalo, e o defeito é silencioso justamente porque não aparece em teste manual, feito devagar.
Antes de mudar qualquer coisa, reproduza o defeito de propósito: force duas transições do mesmo pedido dentro de um único intervalo de polling e confira o que o painel sabe depois.
#!/usr/bin/env bash
# medir-a-perda.sh — prova a perda de transicao no desenho MINIMO, com numero
set -euo pipefail
# 1. Estado inicial: crie um pedido e confirme que ele esta PENDENTE no painel.
aws dynamodb get-item --table-name painel_operacao \
--key '{"PK":{"S":"PEDIDO#demo-001"},"SK":{"S":"ATUAL"}}' \
--query 'Item.status.S' --output text
# 2. Dentro de UM intervalo de polling (aqui, 2 minutos), execute as tres
# transicoes em sequencia rapida.
aws dynamodb update-item --table-name pedidos_cadencia \
--key '{"PK":{"S":"PEDIDO#demo-001"},"SK":{"S":"METADADOS"}}' \
--update-expression 'SET #s = :v' --expression-attribute-names '{"#s":"status"}' \
--expression-attribute-values '{":v":{"S":"SEPARADO"}}'
sleep 5
aws dynamodb update-item --table-name pedidos_cadencia \
--key '{"PK":{"S":"PEDIDO#demo-001"},"SK":{"S":"METADADOS"}}' \
--update-expression 'SET #s = :v' --expression-attribute-names '{"#s":"status"}' \
--expression-attribute-values '{":v":{"S":"CANCELADO"}}'
# 3. Espere o job de polling rodar (ate 2 min) e confira o painel de novo.
sleep 130
aws dynamodb get-item --table-name painel_operacao \
--key '{"PK":{"S":"PEDIDO#demo-001"},"SK":{"S":"ATUAL"}}' \
--query 'Item.status.S' --output text
# Esperado no desenho MINIMO: "CANCELADO", sem nenhum registro de SEPARADO —
# o painel nunca soube que o pedido passou por la. Na Cadencia, medido em
# 07/ago/2026: 0 alertas de reposicao de estoque disparados em 30 casos assim
# ocorridos no mes anterior, todos descobertos so na contagem fisica.
A perda não é atraso — é ausência, e ela não aparece em teste manual
Quem testa clicando devagar nunca reproduz o defeito, porque cada clique cai num intervalo de polling diferente e a função sempre vê a transição. O defeito só aparece sob volume real, quando duas mudanças do MESMO pedido caem dentro da MESMA janela de 2 minutos — e nesse caso não existe consulta, por mais frequente que seja, capaz de recuperar o estado intermediário depois que ele foi sobrescrito. Reduzir o intervalo torna o defeito mais raro; não o elimina.
Arquitetura para produção: a tabela ouvida, não consultada
Cada peça nova abaixo rastreia a uma linha da tabela de requisitos. Se você não conseguir apontar o requisito, a peça é adorno — e este desenho não tem nenhuma.
- → registro do stream: chave + imagem antiga e nova + tipo de evento
- → grava a transição como item novo (SK = sequência do stream)
- → publica só quando OLD=SEPARADO e NEW=CANCELADO
- → lote que esgotou as tentativas configuradas
- → abre o painel de operação
- → lê o histórico de transições e o status atual
- → IteratorAge por invocação, e erros do lote
- Fora da AWS
- Banco de dados
- Compute
- Integração de apps
- Rede e entrega
- Gestão e governança
Não é o desenho anterior com um Lambda a mais: a tabela deixou de ser CONSULTADA e passou a ser OUVIDA. Cada mudança gera um registro com a imagem antiga e a nova, entregue em segundos — e é essa dupla imagem, não a velocidade, que faz o alerta de reposição existir. Percorra os passos: a ordem que o stream garante é mais estreita do que parece, e é exatamente aí que times erram.
- O stream nasce por chave de partição, não pela tabela inteira. Cada partição da tabela escreve seus registros de mudança num shard interno dedicado a ela, e nenhuma outra partição escreve nesse mesmo shard. Não existe um shard único compartilhado pela tabela toda — é esse desenho que faz a garantia de ordem valer POR CHAVE, e não globalmente.
- NEW_AND_OLD_IMAGES é o que faz a transição existir, não só o estado final. Sem a imagem antiga, o registro diria apenas "agora status = CANCELADO". Com as duas, o consumidor compara e sabe que veio de SEPARADO. É a diferença entre saber ONDE o pedido está e saber POR ONDE ele passou — e a segunda pergunta é a que o alerta de estoque precisa responder.
- Quem consulta os shards é o serviço Lambda, não o seu código. O event source mapping faz o polling dos shards por trás das cortinas, agrupa registros em lotes e invoca a função — tipicamente em poucos segundos após a escrita, não minutos. O código de aplicação nunca chama GetRecords diretamente.
- Cada transição vira um item novo; nada sobrescreve o anterior. O painel_operacao passa a ser um log de transições, com a chave de ordenação derivada da sequência do stream. Duas mudanças no mesmo pedido em menos de um minuto viram DOIS itens — é isso que resolve a perda que o polling tinha.
- O alerta de reposição finalmente tem informação para existir. A regra "separado, depois cancelado" só é decidível com a imagem antiga E a nova no mesmo registro. O consumidor publica no tópico só quando as duas condições batem — e o estoque físico deixa de dessincronizar do sistema.
- Ordem é garantida DENTRO do pedido, não ENTRE pedidos. Se dois pedidos DIFERENTES — de lojas diferentes — mudam de status no mesmo segundo, mas moram em partições diferentes, não há garantia sobre qual dos dois eventos o consumidor processa primeiro. Isso não importa aqui, porque cada linha do painel e cada alerta são avaliados por pedido, nunca comparando um pedido com outro.
- Lote com erro vai para a fila, sem travar as outras partições. Se o consumidor falhar ao processar um registro, o event source mapping esgota as tentativas configuradas e manda o lote para a fila de descarte. O processamento das demais chaves de partição continua normalmente, em vez de o shard inteiro travar esperando aquele registro ser resolvido.
A diferença estrutural em relação ao desenho anterior não é "mais um Lambda": é que a tabela de origem PARA de ser consultada por um relógio externo e PASSA a notificar quem está interessado. O painel deixa de ser um espelho sobrescrito e vira um log de transições, porque só um log preserva o que um espelho apaga.
O ajuste que resolve os dois problemas com a mesma peça
NEW_AND_OLD_IMAGES não foi escolhido para acelerar nada — foi escolhido porque a regra de negócio do alerta de estoque exige o ANTES e o DEPOIS no mesmo registro. A latência em segundos é um efeito colateral de trocar consulta por notificação; o que resolve o alerta que nunca disparava é especificamente a imagem dupla.
Como funciona, ponta a ponta
Os nomes dos campos não são jargão: eles são a prova de que a transição foi capturada. Um registro de stream com `OldImage` e `NewImage` diferentes É a transição — não uma inferência sobre ela.
// O que o event source mapping entrega ao consumidor .NET, um registro do
// lote (DynamoDBEvent.Records[i]). NewImage e OldImage vêm no formato de tipo
// do DynamoDB (wrapper {"S": ...}, {"N": ...}), não em JSON "normal".
{
"eventID": "8f105b654fb0c2c4e6a6f6b0f1a2b3c4",
"eventName": "MODIFY",
"eventSource": "aws:dynamodb",
"awsRegion": "us-east-1",
"dynamodb": {
"ApproximateCreationDateTime": 1754571801,
"Keys": {
"PK": { "S": "PEDIDO#b7a1f0b0-2c44-4e9a-9d31-000000009182" },
"SK": { "S": "METADADOS" }
},
"OldImage": {
"status": { "S": "PENDENTE" },
"GSI3PK": { "S": "LOJA#a10ecb4e-6f21-4a02-8b9d-000000000512" },
"GSI3SK": { "S": "2026-08-07T14:03:21Z" }
},
"NewImage": {
"status": { "S": "SEPARADO" }
// GSI3PK/GSI3SK ausentes: o L16 já os removia ao separar — o item some
// da fila de pendentes, e essa ausência também está retratada aqui.
},
// A ordem DENTRO desta partição (este pedido) é garantida pelo número de
// sequência crescente. Entre pedidos de partições diferentes, não é.
"SequenceNumber": "400000000000000499659",
"SizeBytes": 287,
"StreamViewType": "NEW_AND_OLD_IMAGES"
},
"eventSourceARN": "arn:aws:dynamodb:us-east-1:111122223333:table/pedidos_cadencia/stream/2026-08-07T00:00:00.000"
}Por que o SequenceNumber vira parte da chave do painel, e não um timestamp
A entrega do event source mapping é PELO MENOS UMA VEZ: o mesmo registro pode ser reprocessado depois de uma falha parcial de lote. Um timestamp gerado no momento do processamento seria DIFERENTE a cada tentativa, e cada tentativa criaria uma transição nova — duplicando a linha no painel. O SequenceNumber vem PRONTO do stream, é o mesmo em toda tentativa, e faz o `PutItem` reescrever o mesmo item em vez de criar outro.
As decisões, e o que se perde em cada uma
📋 A mesma Cadência do L16 — 900 lojas, 40 mil pedidos/dia — precisando que um painel interno e um alerta de reposição de estoque reajam a mudanças de status em segundos, sem perder nenhuma transição intermediária, com uma equipe de duas pessoas sem capacidade de operar sharding manual.
É nativo da tabela que o L16 já criou — não exige provisionar nem administrar nenhum recurso de streaming separado. A leitura via gatilho do Lambda não tem cobrança adicional de GetRecords, então o custo novo é só a execução do Lambda, proporcional a mudanças reais. E a granularidade de ordem que ele garante — por chave de partição — é exatamente a que os dois casos de uso precisam: o painel e o alerta avaliam cada pedido isoladamente, nunca comparando a ordem entre pedidos diferentes.
Alt: Kinesis Data Streams for DynamoDB (captura para Kinesis) — Dá retenção de até 1 ano em vez de 24 horas, e permite múltiplos consumidores lendo o mesmo stream com checkpoints independentes — sem competir pelos mesmos registros. É a escolha certa quando mais de um sistema precisaria consumir a mesma mudança de forma independente, ou quando 24 horas de retenção é curto demais para o plano de reprocessamento. É o L63.
Alt: Aumentar a frequência do polling (ex.: a cada 10 segundos) — Reduz a JANELA em que a perda pode ocorrer, mas não a elimina: duas mudanças do mesmo pedido dentro de QUALQUER intervalo, por menor que seja, continuam colapsando na mesma consulta. E o custo de leitura cresce linearmente com a frequência, sem nunca chegar à garantia que o requisito pede.
Alt: DMS lendo Aurora/RDS via WAL (CDC relacional) — É a via equivalente para banco relacional, já construída no L19 sobre Aurora/RDS com slot de replicação lógica. Não se aplica aqui: a origem é DynamoDB, que não tem WAL nem binlog — o mecanismo de captura é o stream nativo da própria tabela, não um produto de replicação externo.
Alt: Outbox manual (tabela de eventos escrita na mesma transação de negócio) — Funcionaria, mas reimplementa em código de aplicação o que o stream já oferece de graça pela infraestrutura: todo caminho de escrita da tabela teria que lembrar de gravar também no outbox, e um caminho esquecido produz o mesmo silêncio que o polling produzia — só que por bug de aplicação, não por limite de arquitetura.
| Decisão | Escolha | Alternativas | Motivo | O que se perde |
|---|---|---|---|---|
| Tipo de visualização do stream | NEW_AND_OLD_IMAGES | KEYS_ONLY; NEW_IMAGE | o alerta de estoque exige comparar o antes com o depois no mesmo registro | cada registro fica maior; KEYS_ONLY bastaria se só a CHAVE do item alterado importasse |
| Consumo do stream | event source mapping do Lambda | aplicação chamando GetRecords manualmente | a AWS administra polling de shard, checkpoint e paralelismo | menos controle fino sobre o ritmo exato de leitura de cada shard |
| Retentativa do lote | MaximumRetryAttempts=3 + DLQ | padrão: tentativas infinitas até o registro expirar | um lote genuinamente quebrado não pode bloquear a partição indefinidamente | sem alerta manual na DLQ, o lote fica esquecido lá sem ninguém perceber |
| Escrita do painel | append-only: cada transição é um item novo | sobrescrever o registro do estado atual | é o que resolve o defeito central: nenhuma transição intermediária desaparece | mais itens gravados por pedido, mais unidade de escrita consumida |
| Identidade da transição no painel | SK = SequenceNumber do stream | gerar um novo identificador a cada execução | entrega pelo menos uma vez exige uma chave determinística para não duplicar | nada relevante — é praticamente de graça de implementar |
| Alcance da garantia de ordem | aceitar ordem só por partição (por pedido) | exigir ordem global entre pedidos diferentes | nenhuma regra de negócio deste módulo compara a ordem ENTRE pedidos | se isso um dia for necessário, a resposta é outro desenho — não uma configuração diferente deste mesmo stream |
A dívida que este laboratório cria, e que ele não paga
Se o consumidor ficar fora do ar por mais de 24 horas — deploy quebrado não percebido, incidente prolongado — os registros mais antigos que a janela de retenção já foram descartados, e não existe "recomeçar do zero" depois disso: é preciso reconciliar com uma varredura pontual na tabela de origem, comparando `updated_at` como um polling de emergência faria. Este módulo não constrói essa reconciliação; ele constrói o caminho feliz e documenta o limite.
Construir: ligar o stream, e o log de transições
O stream não é um recurso à parte: é um atributo da tabela que o L16 já criou. Ligá-lo não move dado nenhum — passa a CAPTURAR o que já estava acontecendo.
# tabela.tf — o stream é um ATRIBUTO da tabela do L16, não uma tabela nova
resource "aws_dynamodb_table" "pedidos" {
name = "pedidos_cadencia"
billing_mode = "PAY_PER_REQUEST"
hash_key = "PK"
range_key = "SK"
attribute { name = "PK" type = "S" }
attribute { name = "SK" type = "S" }
# ... GSI1/GSI2/GSI3 continuam como no L16; omitidos aqui por brevidade.
# As duas linhas que este laboratório acrescenta à tabela do L16.
# NEW_AND_OLD_IMAGES e nao NEW_IMAGE: sem a imagem antiga, o consumidor
# nunca saberia que um pedido "cancelado" tinha passado por "separado" —
# e essa comparacao e o motivo deste laboratorio existir.
stream_enabled = true
stream_view_type = "NEW_AND_OLD_IMAGES"
point_in_time_recovery {
enabled = true
}
}
output "stream_arn" {
value = aws_dynamodb_table.pedidos.stream_arn
description = "ARN do stream; muda se a tabela for recriada — o consumidor referencia este output, nunca uma string fixa"
}
LATEST não captura o que já existia antes de ligar o stream
`starting_position = "LATEST"` significa que o consumidor só vê mudanças a partir do momento em que o event source mapping foi criado — o que já estava no ar antes fica de fora. Para este laboratório está certo: o painel começa vazio e se preenche com o uso normal. Se o requisito fosse reconstruir o painel com o histórico inteiro desde sempre, a resposta não seria `TRIM_HORIZON` — o stream só retém 24 horas — seria uma carga inicial por `Scan`, seguida da entrada em regime pelo stream.
# painel.tf — o log de transicoes, e o alerta que passa a ter informacao para existir
resource "aws_dynamodb_table" "painel_operacao" {
name = "painel_operacao"
billing_mode = "PAY_PER_REQUEST"
hash_key = "PK"
range_key = "SK"
# PK = PEDIDO#<id> igual a tabela de origem, para o gerente nao precisar
# aprender uma segunda convencao de chave.
# SK = "ATUAL" -> o estado mais recente conhecido
# SK = "TRANSICAO#<SequenceNumber>" -> uma linha por mudanca, append-only
attribute { name = "PK" type = "S" }
attribute { name = "SK" type = "S" }
tags = { Projeto = "ffv-lab", Papel = "read-model-operacional" }
}
resource "aws_sns_topic" "alerta_estoque" {
name = "ffv-lab-alerta-reposicao-estoque"
}
resource "aws_sqs_queue" "dlq_stream" {
name = "ffv-lab-consumidor-pedidos-dlq"
message_retention_seconds = 1209600 # 14 dias — tempo para alguem investigar antes de expirar
}
Construir: o consumidor, o papel e o event source mapping
Três peças amarradas por comentário: o papel com permissão mínima, a função, e o mapping que liga uma coisa na outra. Nenhuma delas funciona sozinha.
# consumidor.tf — o Lambda que ouve o stream, e o event source mapping que o alimenta
data "aws_iam_policy_document" "consumidor_assume" {
statement {
effect = "Allow"
actions = ["sts:AssumeRole"]
principals {
type = "Service"
identifiers = ["lambda.amazonaws.com"]
}
}
}
resource "aws_iam_role" "consumidor" {
name = "ffv-lab-consumidor-pedidos"
assume_role_policy = data.aws_iam_policy_document.consumidor_assume.json
}
data "aws_iam_policy_document" "consumidor_permissoes" {
# Ler o stream. As tres primeiras aceitam o ARN especifico do stream; a
# quarta e operacao de DESCOBERTA DE CONTA e nao aceita recurso — mesma
# justificativa do ecr:GetAuthorizationToken no L03.
statement {
effect = "Allow"
actions = [
"dynamodb:DescribeStream",
"dynamodb:GetRecords",
"dynamodb:GetShardIterator",
]
resources = ["${aws_dynamodb_table.pedidos.stream_arn}"]
}
statement {
effect = "Allow"
actions = ["dynamodb:ListStreams"]
resources = ["*"] # ListStreams e por conta/regiao; nao existe ARN de stream para restringir
}
statement {
effect = "Allow"
actions = ["dynamodb:PutItem"]
resources = [aws_dynamodb_table.painel_operacao.arn]
}
statement {
effect = "Allow"
actions = ["sns:Publish"]
resources = [aws_sns_topic.alerta_estoque.arn]
}
# O event source mapping usa o papel DA FUNCAO para escrever na DLQ quando
# descarta um lote — por isso a permissao mora aqui, nao num papel separado.
statement {
effect = "Allow"
actions = ["sqs:SendMessage"]
resources = [aws_sqs_queue.dlq_stream.arn]
}
}
resource "aws_iam_role_policy" "consumidor" {
role = aws_iam_role.consumidor.id
policy = data.aws_iam_policy_document.consumidor_permissoes.json
}
resource "aws_lambda_function" "consumidor" {
function_name = "ffv-lab-consumidor-pedidos"
role = aws_iam_role.consumidor.arn
runtime = "dotnet8"
handler = "Consumidor::Consumidor.Function::HandlerAsync"
filename = "consumidor.zip"
timeout = 30
memory_size = 256
environment {
variables = {
TABELA_PAINEL = aws_dynamodb_table.painel_operacao.name
TOPICO_ESTOQUE = aws_sns_topic.alerta_estoque.arn
}
}
}
resource "aws_lambda_event_source_mapping" "stream_pedidos" {
event_source_arn = aws_dynamodb_table.pedidos.stream_arn
function_name = aws_lambda_function.consumidor.arn
starting_position = "LATEST" # so o que acontecer DAQUI PRA FRENTE; ver callout sobre TRIM_HORIZON
batch_size = 100 # padrao do stream do DynamoDB; maximo 10000
maximum_batching_window_in_seconds = 1 # nao espera lote encher; prioriza latencia sobre eficiencia
# Falha isolada nao trava a particao inteira: divide o lote ao meio e tenta
# de novo cada metade separadamente, ate achar o registro problematico.
bisect_batch_on_function_error = true
maximum_retry_attempts = 3 # finito, de proposito — ver callout sobre o padrao ser infinito
# Permite ao consumidor reportar QUAIS registros do lote falharam, em vez de
# o lote inteiro ser tratado como falha ou sucesso.
function_response_types = ["ReportBatchItemFailures"]
destination_config {
on_failure {
destination_arn = aws_sqs_queue.dlq_stream.arn
}
}
}
resource "aws_cloudwatch_metric_alarm" "atraso_do_consumidor" {
alarm_name = "ffv-lab-consumidor-pedidos-iterator-age"
namespace = "AWS/Lambda"
metric_name = "IteratorAge"
statistic = "Maximum"
period = 60
evaluation_periods = 3
threshold = 60000 # 60 s em milissegundos — bem acima do que a operacao normal produz
comparison_operator = "GreaterThanThreshold"
treat_missing_data = "notBreaching"
dimensions = {
FunctionName = aws_lambda_function.consumidor.function_name
}
alarm_actions = [aws_sns_topic.alerta_estoque.arn] # reaproveitado para alerta operacional
}
Tentativa finita, e por que isso é intencional
O padrão do event source mapping para fontes de stream é tentar INDEFINIDAMENTE até o registro expirar do stream — o que soa seguro e não é: um lote com um bug de desserialização travaria a partição por até 24 horas, retentando para sempre um registro que nunca vai processar com sucesso. `maximum_retry_attempts = 3` mais a DLQ troca essa espera cega por um sinal explícito: alguém vê a mensagem na fila e investiga.
Construir: o consumidor em C#, com falha parcial de lote
A parte que mais importa aqui não é a chamada ao DynamoDB — é a resposta que a função devolve quando um registro do lote falha e os outros não.
// Function.cs — compara antes e depois, grava a transicao, publica o alerta
// que o polling nunca conseguia disparar.
using Amazon.DynamoDBv2;
using Amazon.DynamoDBv2.Model;
using Amazon.Lambda.Core;
using Amazon.Lambda.DynamoDBEvents;
using Amazon.SimpleNotificationService;
using Amazon.SimpleNotificationService.Model;
namespace Consumidor;
public class Function
{
private readonly IAmazonDynamoDB _dynamo = new AmazonDynamoDBClient();
private readonly IAmazonSimpleNotificationService _sns = new AmazonSimpleNotificationServiceClient();
private readonly string _tabelaPainel = Environment.GetEnvironmentVariable("TABELA_PAINEL")!;
private readonly string _topicoEstoque = Environment.GetEnvironmentVariable("TOPICO_ESTOQUE")!;
// Retorna DynamoDBBatchResponse com os itens que falharam, para o event
// source mapping reprocessar SO ESSES — e nao o lote inteiro. Sem isso,
// um unico registro com formato inesperado faria o lote inteiro repetir,
// inclusive os que ja tinham sido gravados com sucesso.
public async Task<DynamoDBBatchResponse> HandlerAsync(DynamoDBEvent evento, ILambdaContext contexto)
{
var falhas = new List<DynamoDBBatchResponse.BatchItemFailure>();
foreach (var registro in evento.Records)
{
try
{
await ProcessarAsync(registro, contexto);
}
catch (Exception ex)
{
contexto.Logger.LogError($"Falha no registro {registro.Dynamodb.SequenceNumber}: {ex.Message}");
// O identificador tem de ser o SequenceNumber: e a chave que o
// event source mapping usa para saber ONDE retomar.
falhas.Add(new DynamoDBBatchResponse.BatchItemFailure
{
ItemIdentifier = registro.Dynamodb.SequenceNumber
});
}
}
return new DynamoDBBatchResponse { BatchItemFailures = falhas };
}
private async Task ProcessarAsync(DynamoDBEvent.DynamodbStreamRecord registro, ILambdaContext contexto)
{
// INSERT e pedido novo; REMOVE e exclusao fisica. So MODIFY carrega
// uma transicao de status, que e o que este laboratorio trata.
if (registro.EventName != "MODIFY") return;
var pedidoId = registro.Dynamodb.Keys["PK"].S;
var statusAntigo = registro.Dynamodb.OldImage.TryGetValue("status", out var oldAttr) ? oldAttr.S : null;
var statusNovo = registro.Dynamodb.NewImage.TryGetValue("status", out var newAttr) ? newAttr.S : null;
if (statusAntigo == statusNovo) return; // mudou outro atributo; nao e transicao de status
// Chave deterministica a partir do SequenceNumber: reprocessar o MESMO
// registro do stream (entrega e pelo menos uma vez) grava o MESMO
// item, em vez de duplicar a transicao no painel.
var itemTransicao = new Dictionary<string, AttributeValue>
{
["PK"] = new AttributeValue { S = pedidoId },
["SK"] = new AttributeValue { S = $"TRANSICAO#{registro.Dynamodb.SequenceNumber}" },
["de"] = new AttributeValue { S = statusAntigo ?? "INEXISTENTE" },
["para"] = new AttributeValue { S = statusNovo ?? "DESCONHECIDO" },
["em"] = new AttributeValue { S = DateTimeOffset.UtcNow.ToString("O") },
};
await _dynamo.PutItemAsync(_tabelaPainel, itemTransicao);
// O item ATUAL e sobrescrito de proposito — e o unico que representa
// "agora", e nao precisa de historico.
var itemAtual = new Dictionary<string, AttributeValue>
{
["PK"] = new AttributeValue { S = pedidoId },
["SK"] = new AttributeValue { S = "ATUAL" },
["status"] = new AttributeValue { S = statusNovo ?? "DESCONHECIDO" },
["atualizadoEm"] = new AttributeValue { S = DateTimeOffset.UtcNow.ToString("O") },
};
await _dynamo.PutItemAsync(_tabelaPainel, itemAtual);
// A regra que o polling nunca conseguia avaliar: exige o ANTES e o
// DEPOIS no MESMO registro, e so o stream entrega os dois juntos.
if (statusAntigo == "SEPARADO" && statusNovo == "CANCELADO")
{
await _sns.PublishAsync(new PublishRequest
{
TopicArn = _topicoEstoque,
Subject = "Pedido cancelado apos separado — reposicao de estoque",
Message = $"Pedido {pedidoId} foi cancelado depois de ja estar separado. " +
"O item embalado precisa voltar ao estoque fisico."
});
}
}
}
Sem ReportBatchItemFailures, um registro ruim reprocessa o lote inteiro
O comportamento padrão do event source mapping é tratar o lote como TUDO ou NADA: se a função lançar uma exceção não tratada, TODO o lote é considerado falho e reentregue — inclusive os registros que já tinham sido gravados com sucesso no painel. Sem a escrita idempotente (chave por SequenceNumber) isso duplicaria transições; com ela, é só trabalho refeito. `ReportBatchItemFailures` evita o reprocessamento desnecessário dos registros que já deram certo, mas não substitui a idempotência — as duas proteções resolvem problemas diferentes.
Implantar, e provar que nada se perde
Cinco provas. Nenhuma aceita "o painel parece atualizado" como resultado — cada uma tem um número, uma contagem ou uma ordem esperada.
# provas.sh — cinco medicoes; nenhuma aceita "parece que esta funcionando"
TABELA_ORIGEM=pedidos_cadencia
TABELA_PAINEL=painel_operacao
# ── Prova 1: reacao em segundos, nao minutos ──────────────────────────────
INICIO=$(date +%s%3N)
aws dynamodb update-item --table-name $TABELA_ORIGEM \
--key '{"PK":{"S":"PEDIDO#demo-002"},"SK":{"S":"METADADOS"}}' \
--update-expression 'SET #s = :v' --expression-attribute-names '{"#s":"status"}' \
--expression-attribute-values '{":v":{"S":"SEPARADO"}}'
while true; do
R=$(aws dynamodb get-item --table-name $TABELA_PAINEL \
--key '{"PK":{"S":"PEDIDO#demo-002"},"SK":{"S":"ATUAL"}}' \
--query 'Item.status.S' --output text 2>/dev/null || echo "")
[ "$R" = "SEPARADO" ] && break
sleep 0.5
done
FIM=$(date +%s%3N)
echo "atraso: $(( (FIM - INICIO) ))ms"
# Esperado: abaixo de 10.000 ms na maioria das execucoes. Contra os ate
# 120.000 ms do desenho minimo, e a diferenca de ordem de grandeza que o
# stream compra.
# ── Prova 2: nenhuma transicao intermediaria some ─────────────────────────
aws dynamodb update-item --table-name $TABELA_ORIGEM \
--key '{"PK":{"S":"PEDIDO#demo-003"},"SK":{"S":"METADADOS"}}' \
--update-expression 'SET #s = :v' --expression-attribute-names '{"#s":"status"}' \
--expression-attribute-values '{":v":{"S":"SEPARADO"}}'
aws dynamodb update-item --table-name $TABELA_ORIGEM \
--key '{"PK":{"S":"PEDIDO#demo-003"},"SK":{"S":"METADADOS"}}' \
--update-expression 'SET #s = :v' --expression-attribute-names '{"#s":"status"}' \
--expression-attribute-values '{":v":{"S":"CANCELADO"}}'
sleep 5
aws dynamodb query --table-name $TABELA_PAINEL \
--key-condition-expression 'PK = :p AND begins_with(SK, :t)' \
--expression-attribute-values '{":p":{"S":"PEDIDO#demo-003"},":t":{"S":"TRANSICAO#"}}' \
--query 'length(Items)'
# Esperado: 2 (PENDENTE->SEPARADO e SEPARADO->CANCELADO), nao 1.
# ── Prova 3: ordem preservada por chave de particao ───────────────────────
for i in 1 2 3 4 5; do
aws dynamodb update-item --table-name $TABELA_ORIGEM \
--key '{"PK":{"S":"PEDIDO#demo-004"},"SK":{"S":"METADADOS"}}' \
--update-expression 'SET contador = :v' \
--expression-attribute-values "{\":v\":{\"N\":\"$i\"}}"
done
sleep 5
aws dynamodb query --table-name $TABELA_PAINEL \
--key-condition-expression 'PK = :p AND begins_with(SK, :t)' \
--expression-attribute-values '{":p":{"S":"PEDIDO#demo-004"},":t":{"S":"TRANSICAO#"}}' \
--query 'Items[].SK.S' --output text
# Esperado: os SequenceNumber nos SKs aparecem em ordem CRESCENTE — a mesma
# ordem em que as cinco escritas foram enviadas.
# ── Prova 4: reprocessamento nao duplica (idempotencia) ───────────────────
ANTES=$(aws dynamodb query --table-name $TABELA_PAINEL \
--key-condition-expression 'PK = :p' \
--expression-attribute-values '{":p":{"S":"PEDIDO#demo-003"}}' \
--query 'length(Items)')
aws lambda invoke --function-name ffv-lab-consumidor-pedidos \
--payload file://registro-de-teste.json /tmp/saida.json >/dev/null
aws lambda invoke --function-name ffv-lab-consumidor-pedidos \
--payload file://registro-de-teste.json /tmp/saida.json >/dev/null
DEPOIS=$(aws dynamodb query --table-name $TABELA_PAINEL \
--key-condition-expression 'PK = :p' \
--expression-attribute-values '{":p":{"S":"PEDIDO#demo-003"}}' \
--query 'length(Items)')
echo "antes=$ANTES depois=$DEPOIS"
# Esperado: os dois numeros IGUAIS, mesmo invocando o mesmo registro duas vezes.
# ── Prova 5: o atraso do consumidor, medido pelo IteratorAge ──────────────
aws cloudwatch get-metric-statistics --namespace AWS/Lambda \
--metric-name IteratorAge --statistics Maximum --period 60 \
--start-time "$(date -u -d '10 minutes ago' +%FT%TZ)" --end-time "$(date -u +%FT%TZ)" \
--dimensions Name=FunctionName,Value=ffv-lab-consumidor-pedidos \
--query 'Datapoints[].Maximum'
# Esperado: valores na casa de milissegundos a poucos segundos, em operacao
# normal. Um numero crescendo ao longo de varios pontos indica que o
# consumidor esta processando mais devagar do que a tabela escreve.
| Prova | O que mede | Resultado que aprova | O que reprova, e o que significa |
|---|---|---|---|
| 1 · Latência do consumidor | tempo entre o `UpdateItem` e o item aparecer no painel | abaixo de 10 segundos na maioria das execuções | acima disso, confira `IteratorAge`: o consumidor pode estar processando mais devagar do que a tabela escreve |
| 2 · Nenhuma transição some | contagem de itens `TRANSICAO#` para um pedido com 2 mudanças | exatamente 2, uma por mudança de status | 1 significa que a segunda mudança sobrescreveu a primeira — sinal de que a chave do painel não está usando o SequenceNumber |
| 3 · Ordem preservada por partição | SequenceNumber dos SKs de um mesmo pedido, em sequência | ordem crescente, igual à ordem de envio das escritas | fora de ordem indicaria processamento paralelo cruzando registros da MESMA partição, o que o modelo de consumo não deveria permitir |
| 4 · Reprocessamento não duplica | contagem de itens antes e depois de invocar o mesmo evento duas vezes | os dois números iguais | números diferentes indicam que a chave de escrita não é determinística |
| 5 · Atraso do consumidor sob operação normal | `IteratorAge` via CloudWatch | casa de milissegundos a poucos segundos | crescendo ao longo de vários pontos: o consumidor está acumulando atraso, não só tendo picos pontuais |
Fora do ar por mais de 24 horas: não há "recomeçar de onde parou"
A retenção do stream é fixa em 24 horas — não é ajustável para cima. Se o consumidor ficar parado além desse prazo, os registros mais antigos já foram descartados pela própria infraestrutura, e ligá-lo de volta simplesmente não os encontra mais. A única forma de fechar a lacuna é uma reconciliação pontual — uma varredura na tabela de origem comparando `updated_at` — e ela tem exatamente o mesmo defeito do desenho mínimo deste laboratório: perde transição intermediária que ocorreu durante a lacuna. É um limite real deste desenho, não um detalhe de operação.
Quebrar de propósito: três falhas e o diagnóstico
As três se parecem no sintoma superficial — "o painel ou o alerta não refletem a realidade" — mas cada uma tem uma causa e uma correção diferentes.
| Falha | Como provocar | Sintoma | Onde olhar | Correção |
|---|---|---|---|---|
| Stream com NEW_IMAGE em vez de NEW_AND_OLD_IMAGES | mude `stream_view_type` para `"NEW_IMAGE"` e refaça a prova 2 | o painel atualiza normalmente, mas o alerta de reposição NUNCA dispara — mesmo com o pedido claramente passando por SEPARADO antes de CANCELADO | `describe-table` → `StreamSpecification.StreamViewType` | NEW_AND_OLD_IMAGES; sem a imagem antiga a comparação de transição não existe |
| Chave do painel não determinística | troque `SequenceNumber` por `Guid.NewGuid()` na chave e reprocesse o mesmo lote duas vezes | o número de transições cresce a cada retentativa, mesmo sem nenhuma mudança real no pedido | contagem de itens `TRANSICAO#` crescendo sem novo `UpdateItem` correspondente | chave derivada do `SequenceNumber`, que é o mesmo em toda tentativa do mesmo registro |
| Consumidor desligado por mais de 24 horas | pare o event source mapping, espere além da retenção, religue | um intervalo inteiro de mudanças simplesmente não aparece no painel, sem erro nenhum no CloudWatch | `IteratorAge` não mostra nada de anormal, porque não há mais registro para medir o atraso | reconciliação pontual comparando `updated_at`, aceitando o mesmo risco de perda do desenho mínimo só para o período da lacuna |
Um job compara `updated_at` a cada 2 minutos. Um pedido muda PENDENTE → SEPARADO → CANCELADO em 40 segundos, todos dentro do mesmo intervalo. O que o job vê na próxima consulta?
Segurança: o que passa a trafegar por um caminho novo
O stream entrega a imagem COMPLETA do item, não só o campo que mudou. Se o item tiver atributo sensível, ele agora passa por um caminho novo — o consumidor, seus logs, e a eventual DLQ.
| Risco | Probabilidade | Impacto | Prevenção | Detecção | Resposta |
|---|---|---|---|---|---|
| Papel do consumidor com permissão além do necessário | média | médio | ações restritas ao ARN do stream e da tabela do painel; `ListStreams` é a única exceção justificada, por ser operação de conta | IAM Access Analyzer sobre uso real | reduzir ao uso medido, como no L41 |
| Atributo sensível do pedido vazando em log do consumidor | média | alto | nunca logar `OldImage`/`NewImage` inteiros; logar só os campos relevantes à decisão (neste caso, `status`) | busca por padrão de dado sensível no grupo de logs | redigir o log; se já vazou, tratar como incidente de exposição de dado |
| Mensagem presa na DLQ sem ninguém notando | média | baixo | alarme sobre `ApproximateNumberOfMessages` da fila | CloudWatch no atributo da fila | investigar e reprocessar manualmente, ou descartar com justificativa registrada |
| Consumidor consegue escrever fora da tabela do painel | baixa | médio | `dynamodb:PutItem` restrito ao ARN específico de `painel_operacao` | CloudTrail em `PutItem` fora do ARN esperado | revisar a política; o erro apareceria como `AccessDenied`, não como escrita indevida |
| Tópico de alerta recebendo assinante não autorizado | baixa | médio | política do tópico SNS restrita a quem deveria assinar | auditoria periódica de assinaturas do tópico | remover a assinatura; investigar como foi criada |
O `*` que aparece na política, e por que ele se justifica
`dynamodb:ListStreams` é uma operação de CONTA/REGIÃO — ela lista os streams existentes antes de você saber o ARN de nenhum, e por isso não aceita recurso específico. As demais ações desta política são restritas ao ARN exato do stream ou da tabela. A regra não é "nunca use `*`" — é "todo `*` tem de vir com a frase que explica por que não pode ser mais estreito".
Observabilidade: as perguntas que o painel de operação do CONSUMIDOR responde
Um consumidor de stream tem uma pergunta central: ele está acompanhando o ritmo da tabela, ou está ficando para trás? Todo o resto decorre disso.
| Pergunta | Métrica ou consulta | O que significa mudar | Limiar inicial |
|---|---|---|---|
| O consumidor está atrasado em relação à escrita? | `IteratorAge` (AWS/Lambda) | tempo entre o registro ser escrito no stream e ser entregue à função | > 60 s por 3 períodos de 1 min |
| Há lote falhando sistematicamente? | `Errors` da função associada ao event source mapping | erro recorrente indica bug de processamento, não falha pontual | > 0 sustentado |
| Mensagens estão se acumulando na fila de descarte? | `ApproximateNumberOfMessages` da DLQ | lote que ninguém está reprocessando | > 0 por mais de 1 hora |
| Quantas transições estão sendo gravadas por minuto? | contagem de `PutItem` no painel | queda abrupta pode indicar consumidor parado, não ausência real de mudança | comparar com a linha de base |
| O alerta de estoque está disparando na frequência esperada? | contagem de publicações no tópico SNS | zero por muito tempo pode ser bom sinal ou consumidor silenciosamente quebrado | cruzar com o volume de cancelamentos, não olhar isolado |
A métrica que some exatamente quando mais precisa dela
Sob um pico de escrita — uma campanha, por exemplo — `IteratorAge` é justamente a métrica que sobe, e é justamente quando a atenção do time está no painel de negócio, não no de infraestrutura. Alarme automático nesta métrica substitui a vigilância manual que a equipe pequena não tem tempo de fazer durante o próprio pico que gerou o atraso.
Escala: 10, 10 mil, 1 milhão de mudanças por dia
| Volume | O que acontece com o consumo | O que passa a doer | O que fazer |
|---|---|---|---|
| 10 mudanças/dia | um shard, invocações esparsas | nada; é o cenário do laboratório | nada |
| 10 mil mudanças/dia (900 lojas, regime normal) | poucos shards ativos, lotes pequenos | nenhum ajuste especial; `maximum_batching_window` de 1 s prioriza latência | confirmar que `IteratorAge` fica na casa de segundos sob esse regime |
| 1 milhão de mudanças/dia (campanha, 20× o normal) | mais partições escrevendo ao mesmo tempo, mais shards, Lambda escalando concorrência automaticamente | concorrência do Lambda pode esbarrar em cota de conta se outras funções competem pelo mesmo limite | verificar cota de concorrência reservada; considerar `parallelization_factor` acima de 1 por shard se um shard concentrar volume |
| Falha de AZ | DynamoDB, Streams e Lambda são serviços gerenciados multi-AZ por padrão | nenhuma ação manual — não há instância nem shard fixo para religar | nada a fazer aqui; é diferente do ALB/ECS do L01, onde AZ é decisão explícita de rede |
| Uma loja concentra volume muito acima das outras | a partição dela concentra escrita e leitura do stream | possível hot partition — o mesmo risco que o L16 já registrou como dívida aberta | sufixo de dispersão na chave, quando o volume real justificar |
O que NÃO cresce com o volume: o custo por transição
Diferente do Scan do desenho errado do L16, o custo do consumo por stream é proporcional a QUANTAS mudanças acontecem, nunca a quantos itens a tabela tem no total. Uma tabela com 40 milhões de pedidos acumulados e uma com 40 mil pagam o mesmo por transição — a diferença é só o volume de transições em si.
Concorrência do Lambda é um limite de CONTA, não só desta função
Sob um pico de 20×, o consumidor pode escalar concorrência o suficiente para esbarrar na cota da conta — que é compartilhada com toda função Lambda daquela conta, não só a deste laboratório. O sintoma é sutil: outra função, sem relação nenhuma com pedidos, começa a ser limitada (throttled) durante a campanha, porque o consumidor do stream consumiu a fatia disponível. Concorrência reservada para o consumidor evita que ele roube capacidade de outra função, mas também limita o quanto ELE pode escalar — é uma troca, não uma folga de graça.
Custo: o que este laboratório acrescenta à fatura
O stream em si não tem cobrança própria, e a leitura via gatilho do Lambda para DynamoDB Streams não é cobrada separadamente — o que se paga é a execução da função, proporcional a mudanças reais, e a tabela nova do painel.
| Cenário | Volume | O que acrescenta | Tendência | Otimização |
|---|---|---|---|---|
| Protótipo | algumas dezenas de mudanças/dia | invocações esporádicas do Lambda; tabela do painel quase vazia | desprezível | nenhuma; otimizar aqui é gastar atenção onde não há dinheiro |
| Produção atual | ~40 mil pedidos/dia, 2 a 4 mudanças cada | dezenas de milhares de invocações/dia, cada uma processando um lote pequeno | baixa e previsível, cresce com o número de PEDIDOS, não com o tamanho da tabela | lote maior (`batch_size`) reduz número de invocações às custas de um pouco de latência |
| Pico de campanha | 20× o volume normal, em rajadas | mais concorrência simultânea do Lambda; mais unidade de escrita no painel | sob demanda no painel absorve sem replanejamento | confirmar cota de concorrência reservada antes do pico, não durante |
| Dimensão | Cobra por | Cuidado |
|---|---|---|
| Lambda (consumidor) | GB-segundo de execução, por invocação | proporcional a mudanças reais — é a inversão do defeito do Scan, que crescia com o tamanho da tabela |
| DynamoDB Streams | nada além da tabela de origem, para leitura via gatilho do Lambda | consumidores QUE NÃO são gatilho do Lambda pagam por leitura acima de uma cota mensal gratuita — não é o caso deste desenho |
| Tabela painel_operacao | unidade de escrita/leitura sob demanda | append-only grava mais itens por pedido que um espelho sobrescrito — mais WCU que o desenho mínimo, e é o preço de não perder transição |
| SQS (DLQ) | por requisição, e por GB-mês de retenção | valor pequeno em operação normal; sobe se o consumidor começar a falhar sistematicamente |
| SNS | por publicação e por notificação entregue | proporcional a quantos alertas de fato disparam — baixo por natureza |
| CloudWatch | métrica customizada e alarme-mês | valor pequeno e fixo; não é onde se economiza |
O ganho de custo que não está em nenhuma linha da AWS
Os 30 casos mensais de item embalado e não devolvido ao estoque, encontrados só na auditoria física, representavam perda de inventário que nenhuma fatura da AWS jamais mostraria. O que este laboratório resolve não aparece como redução de custo de nuvem — aparece como redução de perda operacional, que estava sendo descoberta tarde demais para ser corrigida.
Well-Architected nos seis pilares
| Pilar | Situação ao fim deste laboratório | Risco que fica | Melhoria | Prioridade |
|---|---|---|---|---|
| Excelência operacional | consumidor reagindo em segundos, com falha parcial de lote isolada e observável | sem reconciliação automática para downtime acima de 24 h | rotina de reconciliação por varredura, disparada por alarme de gap detectado | média |
| Segurança | permissão restrita por ARN de stream e de tabela; log sem imagem completa do item | DLQ sem verificação automática de conteúdo sensível antes de reter 14 dias | política de retenção mais curta na DLQ, ou redação automática antes de gravar | média |
| Confiabilidade | ordem preservada por chave, escrita idempotente, bisecção de lote em erro | ordem NÃO garantida entre pedidos diferentes — aceitável hoje, documentado como limite | se algum dia precisar de ordem global, é outro desenho — não uma configuração deste | baixa |
| Eficiência de performance | latência em segundos, medida por `IteratorAge` | lote pequeno com `maximum_batching_window=1s` prioriza latência sobre eficiência de custo | aumentar a janela de batching se o custo por invocação começar a importar mais que a latência | baixa |
| Otimização de custos | custo proporcional a mudanças reais, não ao tamanho da tabela | nenhum hoje, no volume atual | revisar `batch_size` se o volume crescer uma ordem de grandeza | baixa |
| Sustentabilidade | nenhum recurso ocioso: sem servidor fixo, sem shard provisionado sem uso | DLQ retendo mensagem por 14 dias mesmo se ninguém for olhar | alarme na DLQ para forçar decisão em vez de deixar reter até expirar | baixa |
Evolução em níveis: o que muda, e o que passa a doer
A terceira arquitetura não é um desenho: é a resposta a QUANDO trocar de mecanismo de notificação. Cada nível resolve um risco e compra outro.
Painel sem sincronização nenhuma: o gerente pergunta ao suporte por telefone. É de onde a Cadência partiu antes mesmo do polling.Função agendada comparando `updated_at`. É o desenho mínimo deste laboratório, e continua aceitável quando duas mudanças no mesmo item dentro do intervalo são raras e toleráveis.DynamoDB Streams com NEW_AND_OLD_IMAGES, consumidor via event source mapping, escrita idempotente e append-only.Kinesis Data Streams for DynamoDB: retenção de até 1 ano, fan-out para vários sistemas lendo o mesmo stream com checkpoints próprios (L63).O consumidor republica a transição num barramento (EventBridge) com contrato de evento versionado entre times, e replay controlado (L24).O log de transições, acumulado ao longo de meses, vira base para detectar padrão anormal de SEPARADO→CANCELADO por loja — candidato a erro de separação ou fraude interna.A ordem não é negociável, e o motivo é concreto
O nível 6 depende de um histórico de transições confiável, que depende de nenhuma delas ter se perdido — que é exatamente o que o nível 3 resolve. Treinar um classificador sobre dado que o nível 2 produziria seria treinar sobre um histórico com buracos que ninguém sabe onde estão, porque o próprio mecanismo de coleta os escondia.
Onde IA entra nesta arquitetura, e onde não entra
Neste módulo, IA não resolve o problema central, e forçá-la seria o antipadrão que a própria série critica. "Como reagir a uma mudança em segundos, sem perder nenhuma e com ordem preservada" tem resposta determinística: CDC nativo, imagem dupla, chave idempotente. Um modelo não melhora nenhuma dessas três — são mecanismo de infraestrutura, não julgamento.
Há um lugar onde IA agregaria valor real, e ele só existe DEPOIS deste laboratório: com o histórico de transições fluindo de forma confiável, é possível treinar um classificador que sinaliza um padrão anormal — por exemplo, uma loja específica com taxa de SEPARADO→CANCELADO muito acima da média das outras, candidato a erro de processo ou fraude interna, antes que o volume vire prejuízo relevante.
| Pergunta | Resposta honesta para este módulo |
|---|---|
| Qual problema a IA resolveria? | detectar um padrão anormal na taxa de transição SEPARADO→CANCELADO por loja, cedo o suficiente para investigar antes do prejuízo somar |
| Por que uma regra não bastaria? | uma regra de limiar fixo bastaria para começar — "mais de N cancelamentos pós-separação por dia" cobre o caso óbvio. IA só se justifica quando o padrão normal varia por loja, sazonalidade ou porte, e um limiar único gera alarme falso demais para ser útil |
| De onde viriam os dados? | o próprio log de transições que este laboratório grava — nenhum dado novo a coletar |
| Qual o risco? | aprender de poucas lojas com histórico longo e reprovar loja pequena com amostra insuficiente; exige avaliação e um caminho de revisão humana sempre disponível, nunca decisão automática de bloqueio |
| Por que não agora? | porque este laboratório acabou de nascer: sem meses de histórico confiável, treinar qualquer modelo aqui seria superstição com aparência de estatística |
O uso de IA que parece atraente e é armadilha aqui
Pedir a um modelo para "olhar o evento do stream e decidir se deve alertar o estoque" trocaria uma regra determinística — OLD=SEPARADO e NEW=CANCELADO, exata e auditável — por uma decisão probabilística sobre algo que já é 100% observável no próprio registro. Onde existe comparação direta de dois valores, um modelo só acrescenta latência e uma chance de errar com confiança. O lugar certo para IA é a CAMADA ACIMA, sobre o padrão agregado — não sobre a decisão individual que já tem resposta exata.
Anti-padrões deste laboratório
| Anti-padrão | Por que alguém faz | Por que é problema | Sintoma em produção | Forma correta | Quando é aceitável |
|---|---|---|---|---|---|
| Reduzir o intervalo de polling para "resolver" o atraso | é a menor mudança possível sobre um código que já existe | reduz a JANELA de perda, não a elimina; e multiplica o custo de leitura | perda de transição continua ocorrendo, só que com menos frequência — mais difícil de reproduzir e diagnosticar | trocar consulta por stream | nunca, quando a regra de negócio depende de transição intermediária |
| NEW_IMAGE sozinho, sem OLD_IMAGE | parece suficiente porque "eu só preciso saber o estado atual" | qualquer regra que dependa do ANTES fica impossível de implementar depois — é preciso religar o stream e perder o que já passou | alerta ou lógica que depende de transição específica nunca funciona, sem erro nenhum no código | NEW_AND_OLD_IMAGES sempre que houver qualquer chance de precisar da transição, não só do estado | quando genuinamente só o estado atual importa, e isso está escrito no requisito |
| Timestamp de aplicação como chave de idempotência | parece mais legível que um número de sequência opaco | muda a cada tentativa de reprocessamento, então não protege contra duplicação nenhuma | contagem de itens cresce sem novo evento de negócio correspondente | SequenceNumber do próprio registro do stream | nunca; ele existe e é feito exatamente para isso |
| Tentativas infinitas sem DLQ | é o padrão, e "resolver sozinho" parece mais robusto | um lote genuinamente quebrado trava a partição até o registro expirar do stream, até 24 horas | partição específica atrasada enquanto as outras processam normalmente, sem erro visível | limite finito de tentativas mais DLQ com alarme | nunca em produção; aceitável só em prototipagem isolada |
| Tratar o stream como fila durável para sempre | o nome "stream" sugere durabilidade, e o conceito de fila é familiar | a retenção é fixa em 24 horas; downtime maior perde dado de forma irrecuperável pelo próprio stream | "sumiu um pedaço do histórico e não sabemos por quê" depois de um incidente prolongado | plano de reconciliação explícito para downtime acima da retenção | nunca — é sempre um limite a respeitar, não uma suposição a fazer |
Quando algo não funciona
| Sintoma | Causa provável | Como investigar | Onde olhar | Correção |
|---|---|---|---|---|
| Painel atualiza, mas alerta de estoque nunca dispara | stream sem NEW_AND_OLD_IMAGES, ou comparação de status ausente no código | confira o tipo de visualização e releia a condição no consumidor | `describe-table` → `StreamViewType`; log do consumidor | NEW_AND_OLD_IMAGES; comparar OldImage.status com NewImage.status explicitamente |
| Contagem de transições maior que o número de mudanças reais | chave de escrita não determinística — duplicando em reprocessamento | compare o número de `UpdateItem` na origem com o número de itens `TRANSICAO#` | como a chave `SK` do painel é gerada no código do consumidor | usar o SequenceNumber do registro do stream, nunca timestamp ou GUID gerado na hora |
| IteratorAge crescendo continuamente | consumidor processando mais devagar do que a tabela escreve | compare a taxa de invocação com a taxa de escrita na tabela de origem | `IteratorAge` e `Duration` da função no CloudWatch | aumentar `parallelization_factor` por shard, ou reduzir trabalho por invocação |
| Partição específica sempre atrasada, as outras normais | um registro daquela partição está travando o lote repetidamente | veja se `bisect_batch_on_function_error` está ativo, e confira a DLQ | mensagem correspondente na fila de descarte | corrigir o registro problemático ou o parsing que falha nele; reprocessar manualmente a partir da DLQ |
| Lacuna de dados depois de um incidente prolongado | consumidor ficou fora do ar por mais de 24 horas | compare o horário do incidente com a janela de retenção do stream | histórico de eventos do event source mapping; horário do incidente | reconciliação pontual por varredura, aceitando que a ordem de transição da lacuna se perdeu |
| Função nunca é invocada, mesmo com mudanças na tabela | event source mapping desabilitado, ou papel sem permissão de leitura no stream | confira o estado do mapping e teste a permissão isoladamente | `aws lambda get-event-source-mapping` → `State` | reabilitar o mapping; corrigir a política do papel de execução |
A pergunta que resolve metade destes casos
Antes de mexer em parâmetro, pergunte: o problema é de CAPTURA (a mudança não virou registro no stream) ou de CONSUMO (o registro existe, mas o consumidor não o processou a tempo ou corretamente)? A primeira família aponta para a configuração da tabela e do tipo de visualização; a segunda, para o event source mapping e o código do consumidor. As duas produzem o mesmo sintoma visível — "o painel está errado" — e apontam para metades diferentes do desenho.
Limpeza: o que o destroy não leva
O stream em si morre com a tabela — não é um recurso separado a apagar. O que sobrevive ao terraform destroy é o que fica de pé se a ordem de remoção for invertida, ou o que foi criado fora do estado do Terraform.
# limpar.sh — o que o destroy nao leva
# 1. Derrube o que o Terraform administra.
terraform destroy -auto-approve
# 2. FILA DE DESCARTE COM MENSAGEM: se algum lote foi descartado durante o
# laboratorio, a fila pode nao esvaziar sozinha e continuar cobrando
# centavos por retencao de 14 dias.
aws sqs get-queue-attributes --queue-url "$(terraform output -raw dlq_url 2>/dev/null)" \
--attribute-names ApproximateNumberOfMessages 2>/dev/null || true
# 3. O STREAM NAO E UM RECURSO SEPARADO: ele morre com a tabela. Mas se voce
# manteve a tabela do L16 e so removeu o consumidor, o stream CONTINUA
# LIGADO e continua gerando registros que ninguem le — sem custo de
# armazenamento (streams nao cobram parado), mas sem utilidade nenhuma.
aws dynamodb describe-table --table-name pedidos_cadencia \
--query 'Table.StreamSpecification'
# 4. Confirme que nao sobrou nada com a tag do projeto.
aws resourcegroupstaggingapi get-resources \
--tag-filters Key=Projeto,Values=ffv-lab \
--query "ResourceTagMappingList[].ResourceARN" --output table
| Recurso | Sai no destroy? | Cobra parado? | Por que fica |
|---|---|---|---|
| Stream da tabela | sim, junto com a tabela | não | não é um recurso independente; é um atributo. Mas se você manteve a tabela do L16 e só apagou o consumidor, o stream CONTINUA LIGADO gerando registros que ninguém lê |
| Tabela painel_operacao | sim | sim, por unidade sob demanda enquanto existir | igual a qualquer tabela DynamoDB — cobra por leitura/escrita, não por hora ligada |
| Fila DLQ com mensagem | sim, mas com mensagem dentro se não foi esvaziada antes | sim, retenção de 14 dias | mensagem de lote descartado não é apagada automaticamente pelo destroy |
| Tópico SNS | sim | centavos por publicação, não por estar parado | sem custo relevante de estar ocioso |
| Papel e política IAM | sim | não | nenhum custo, mas confira se não sobrou papel órfão referenciado em outro lugar |
| Alarme CloudWatch | sim se em Terraform | centavos | alarme criado à mão no console não aparece no estado |
Stream ligado sem consumidor não custa, mas confunde quem investigar depois
DynamoDB Streams não é cobrado por ficar parado nem por gerar registros que ninguém lê — mas um stream ligado numa tabela sem event source mapping ativo é exatamente o tipo de configuração órfã que confunde uma investigação futura: alguém encontra o atributo ligado, presume que existe um consumidor, e perde tempo procurando um Lambda que já foi apagado. Desligue `stream_enabled` junto com o consumidor, não separado.
Resumo: problema, peça e motivo
| Problema | Peça | Por que ela, e não outra |
|---|---|---|
| Painel atrasado até 2 minutos | DynamoDB Streams + event source mapping do Lambda | reação em segundos, sem precisar de shard próprio para administrar |
| Transição intermediária desaparecendo | NEW_AND_OLD_IMAGES + escrita append-only | a imagem dupla é o que permite comparar antes/depois; append-only é o que preserva cada mudança como um item |
| Alerta de estoque que nunca disparava | comparação OLD.status=SEPARADO e NEW.status=CANCELADO | só é decidível com as duas imagens no mesmo registro |
| Duplicação em reprocessamento | chave do painel derivada do SequenceNumber | entrega pelo menos uma vez exige uma chave determinística para não duplicar |
| Lote quebrado travando uma partição | tentativas finitas + DLQ + bisecção de lote | isola o registro problemático sem parar o processamento das outras chaves |
| Ordem confundida com garantia global | aceitar ordem só por chave de partição | é a granularidade real que o serviço garante — fingir mais do que isso quebra em produção, não em teste |
| Falha | O que a protege | O que ela NÃO protege |
|---|---|---|
| Transição intermediária perdida | NEW_AND_OLD_IMAGES + append-only | downtime do consumidor acima de 24 horas — aí a perda volta a existir |
| Duplicação por reprocessamento | chave idempotente por SequenceNumber | duplicação por bug de lógica que grava duas transições distintas para o mesmo evento |
| Lote quebrado travando a partição | tentativas finitas + DLQ | a causa raiz do registro quebrado — a DLQ isola, não corrige |
| Atraso não percebido | alarme em `IteratorAge` | atraso que ainda está dentro do limiar mas já maior que o histórico normal — vale olhar tendência, não só limiar fixo |
| Ordem trocada dentro do mesmo pedido | shard dedicado por partição | ordem relativa entre pedidos DIFERENTES — nunca foi garantida |
- Expedidor muda o status do pedido com um UpdateItem comum, igual ao L16.
- A tabela, com o stream ligado, grava um registro com a imagem antes e depois no shard da partição.
- O event source mapping do Lambda entrega o registro num lote, em poucos segundos.
- O consumidor distingue INSERT/MODIFY/REMOVE, e só processa MODIFY com status alterado.
- Grava a transição no painel com SK derivado do SequenceNumber — idempotente por construção.
- Se OLD=SEPARADO e NEW=CANCELADO, publica o alerta de reposição de estoque.
- Se o processamento de um registro falha, o event source mapping isola o lote e tenta de novo.
- Após as tentativas configuradas, o lote vai para a DLQ, sem travar as outras partições.
- O gerente abre o painel e vê o histórico completo, com segundos de atraso.
- Uma segunda transição do mesmo pedido, minutos depois, aparece como um SEGUNDO item — nada some.
Perguntas frequentes
❓ Por que um job que consulta updated_at a cada minuto ainda perde mudança de status?
❓ DynamoDB Streams garante que os eventos chegam na ordem em que aconteceram?
❓ Por quanto tempo o DynamoDB Streams guarda os registros de mudança?
❓ Preciso da imagem antiga (OLD) se só quero saber o estado atual de um item?
❓ DynamoDB Streams garante exactly-once na entrega para o Lambda?
❓ Qual a diferença entre DynamoDB Streams e Kinesis Data Streams for DynamoDB?
❓ O que acontece se meu consumidor ficar fora do ar por mais de 24 horas?
❓ Preciso pagar por GetRecords quando uso DynamoDB Streams com Lambda?
Fixando
Os pedidos P1 (loja A) e P2 (loja B) mudam de status no mesmo segundo. O que o DynamoDB Streams garante sobre a ordem relativa entre esses dois eventos?
O consumidor grava uma transição no painel usando o `SequenceNumber` do registro do stream como parte da chave, em vez de um timestamp gerado pela aplicação no momento do processamento. Por quê?
Conhecimentos, próximo módulo e documentação
| Item | Conteúdo |
|---|---|
| Conhecimentos anteriores necessários | L16 no ar (tabela pedidos_cadencia com GSIs), .NET 8, Docker, Git e Terraform básicos |
| Conhecimentos adquiridos | por que polling por timestamp perde transição, não só atrasa; NEW_AND_OLD_IMAGES como requisito de regras que dependem do antes e do depois; ordem garantida por partição, não globalmente; retenção de 24 h como limite real; idempotência via SequenceNumber diante de entrega pelo menos uma vez |
| Limitação que fica | sem reconciliação automática para downtime acima de 24 h; ordem entre pedidos diferentes nunca foi garantida e continua não sendo |
| Próximo exemplo recomendado | L63 — ingestão em streaming com Kinesis: shard, chave de partição e reprocesso a partir do início, para quando um único stream de 24 h deixa de ser suficiente |
| Também habilitado por este módulo | qualquer sistema secundário que hoje consulta pedidos_cadencia periodicamente pode migrar para o mesmo consumidor, ou um novo, lendo o mesmo stream |
| Data da última validação técnica | 7 de agosto de 2026 |
Documentação oficial consultada: Change data capture for DynamoDB Streams — retenção de 24 horas, tipos de visualização e a garantia de ordem por partição; Process DynamoDB records with Lambda e Configuring partial batch response with DynamoDB and Lambda — tamanho de lote, tentativas, destino de falha e o formato de resposta do `ReportBatchItemFailures`; e a página de preços do DynamoDB, quanto à isenção de cobrança quando o consumo é feito via gatilho do Lambda. Os valores de preço não aparecem neste módulo por decisão: use o AWS Pricing Calculator, porque preço varia por região e envelhece mais rápido que o conteúdo.
O que não foi verificado, e você deve conferir na sua conta
Os números de latência citados — abaixo de 10 segundos na maioria das execuções — vêm de uma medição de exemplo em volume baixo, e servem como ordem de grandeza, não como referência. Sob pico real, `IteratorAge` pode subir por minutos até o Lambda escalar concorrência o suficiente; meça na SUA conta, sob a SUA carga de pico, antes de definir o limiar do alarme. Da mesma forma, os "30 casos mensais" de reposição perdida são uma hipótese narrativa deste laboratório, não um dado da AWS.
Terminou de ler?
Marcar como concluído registra o XP, mantém sua sequência e coloca 3 cartas deste módulo na fila de revisão espaçada.
Próximos passos sugeridos
Temas deste módulo
Discussão
Carregando comentários…