7.6 Dados em tempo real e streaming
Visão geral e motivação
Quase tudo que vocês sabem sobre pipelines de dados pressupõe que os dados ficam parados. Vocês coletam um dia de registros, rodam um job durante a noite e leem os resultados pela manhã. Os dados em tempo real e em streaming invertem essa suposição. Em vez de processar uma pilha de dados já pronta, vocês processam um fluxo sem fim de eventos à medida que chegam e produzem respostas continuamente. Essa é a diferença entre o processamento em lote (batch), que opera sobre um conjunto de dados delimitado e completo, e o processamento de fluxo, que opera sobre um fluxo ilimitado e nunca terminado.
Para grandes equipes, o streaming aparece no momento em que a latência passa a importar para o negócio. Uma decisão de fraude que chega uma hora atrasada não vale nada. Um sinal de personalização que aterrissa amanhã não personaliza nada. Um painel operacional que atrasa em relação à realidade por um turno engana quem o observa. O capítulo 7.2 (engenharia de dados) argumenta que vocês devem escolher o batch por padrão e recorrer ao streaming apenas onde a latência realmente compensa, e este capítulo leva vocês o resto do caminho: quando o tempo real merece o seu custo e como construí-lo sem pôr fogo no orçamento operacional. O streaming fica perto dos padrões de mensagens orientadas a eventos do capítulo 3.12 (arquitetura orientada a eventos e mensageria), das escolhas de armazenamento do capítulo 3.4 (arquitetura e armazenamento de dados) e das práticas de telemetria do capítulo 9.2 (observabilidade e telemetria).
Os contextos corporativo e governamental elevam as apostas. Um banco pontua transações quanto à fraude no tempo que um leitor de cartão leva para piscar. Uma agência de transporte rastreia veículos e prevê chegadas para milhões de passageiros. Uma agência de benefícios vigia anomalias em pedidos enquanto mantém um registro auditável de cada decisão. Em todos esses casos, o valor vem de agir sobre os dados enquanto ainda estão frescos, e o risco vem de agir sobre dados errados, incompletos ou impossíveis de reconstruir depois. Este capítulo tem opinião firme sobre ambos.
Princípios fundamentais
- Recorram ao streaming apenas quando a latência tem valor de negócio claro. O batch é mais barato e mais simples.
- Distingam dados delimitados (finitos) de dados ilimitados (sem fim) e projetem de acordo.
- Tratem o tempo do evento, não o de chegada, como fonte da verdade e planejem para dados atrasados e fora de ordem.
- Janelas e marcas d’água são como vocês obtêm respostas finitas de fluxos infinitos.
- Prefiram resultados efetivamente-uma-vez por meio de destinos idempotentes a promessas frágeis de exatamente-uma-vez.
- O processamento com estado precisa de checkpoints para se recuperar sem perder nem contar em dobro.
- Projetem para a contrapressão e o reprocessamento desde o primeiro dia, não como uma reflexão tardia.
- Mantenham a lógica de streaming observável e auditável. Um fluxo silencioso é pior que um batch que falhou.
Recomendações
Justifique o tempo real antes de construí-lo
A decisão de streaming mais importante é se vão fazer streaming. O tempo real aproximadamente dobra a complexidade operacional e o custo, porque vocês trocam um job que roda e para por um sistema que precisa se manter saudável a cada segundo. Antes de se comprometerem, nomeiem a decisão que os dados frescos viabilizam e o custo de essa decisão chegar tarde. A pontuação de fraude, o alerta operacional e a personalização ao vivo costumam passar do patamar. Um painel que um humano olha duas vezes por dia quase nunca passa, por mais satisfatório que “tempo real” soe numa reunião de planejamento. Escrevam o requisito de latência como um número, em segundos ou minutos, e confrontem-no com a realidade. Muito do que as pessoas chamam de tempo real é bem atendido por micro-lotes que rodam a cada poucos minutos por uma fração do custo.
Projete em torno do tempo do evento, não do tempo de processamento
A ideia mais difícil do streaming é que os eventos acontecem num momento e são processados em outro. O tempo do evento é quando a coisa de fato ocorreu, por exemplo quando um passageiro encostou o cartão. O tempo de processamento é quando o seu sistema chegou a tratá-la. Os dois se afastam constantemente: um telefone perde o sinal num túnel e envia três minutos de toques de uma vez, um soluço de rede reordena mensagens, uma partição fica atrasada. Se vocês calculam sobre o tempo de processamento, os seus números oscilam com a infraestrutura em vez de refletir o mundo. Esse problema de dados atrasados e fora de ordem é o coração da disciplina e se liga diretamente à modelagem de eventos da arquitetura orientada a eventos. Carimbem cada evento com seu tempo de evento na origem, carreguem esse carimbo por todo o pipeline e calculem os resultados contra ele.
Use janelas e marcas d’água para obter respostas finitas
Um fluxo ilimitado nunca acaba, então “contar os eventos” não tem resposta até vocês delimitá-lo. As janelas fazem essa delimitação. As janelas fixas (tumbling) cortam o tempo em baldes fixos e sem sobreposição, por exemplo a cada minuto. As janelas deslizantes se sobrepõem, de modo que uma janela de cinco minutos que avança a cada minuto dá uma cifra móvel suave. As janelas de sessão agrupam rajadas de atividade separadas por intervalos de inatividade, o que serve bem às sessões de usuário. Tendo janelas, vocês precisam decidir quando uma janela está concluída, porque dados atrasados ainda podem chegar. Uma marca d’água (watermark) é a estimativa do sistema de que provavelmente já viu todos os eventos até um dado tempo de evento. Quando a marca d’água passa do fim de uma janela, vocês emitem o resultado. Ajustem quanto esperam: mantendo as janelas abertas por mais tempo, toleram mais atraso ao custo de latência e memória; fechando-as mais depressa, arriscam descartar retardatários. Decidam explicitamente o que acontece com os dados que chegam depois de uma janela fechar, se vocês os descartam, os registram ou emitem uma correção.
Torne os destinos idempotentes e prefira o efetivamente-uma-vez
As garantias de entrega parecem simples e não são. A entrega pelo-menos-uma-vez significa que todo evento é processado, mas alguns podem ser processados mais de uma vez após uma nova tentativa, de modo que as contagens podem inflar. O exatamente-uma-vez soa ideal mas é caro e, tomado literalmente entre sistemas externos arbitrários, muitas vezes impossível. A meta prática é o efetivamente-uma-vez: o resultado observável é como se cada evento tivesse sido processado uma só vez, mesmo que a maquinaria tenha tentado de novo por baixo. Vocês chegam lá tornando os destinos idempotentes, seguros para escrita repetida, usando chaves determinísticas e upserts para que um evento reproduzido sobrescreva em vez de duplicar. Combinem a entrega pelo-menos-uma-vez com escritas idempotentes e vocês obtêm resultados corretos sem pagar por coordenação transacional pesada em toda parte. Reservem a maquinaria de exatamente-uma-vez de verdade para os poucos lugares que genuinamente a exigem.
Faça checkpoints do processamento com estado para que ele se recupere
Muitas computações úteis de streaming têm estado: contagens correntes, junções entre fluxos, desduplicação, modelos de fraude que lembram o comportamento recente. Esse estado vive na memória e sumiria quando um processo reiniciasse. O checkpoint tira periodicamente um instantâneo do estado e da posição no fluxo juntos, de modo que, após uma queda, o sistema retoma de um ponto consistente em vez de reproduzir tudo ou perder a memória. Dimensionem o estado deliberadamente, porque o estado ilimitado é um jeito comum de esgotar a memória de um job de streaming em produção. Usem expiração e tempo de vida no estado de que não precisam mais e monitorem o tamanho do estado como uma métrica de primeira classe. O tempo de recuperação após uma falha é uma preocupação real de nível de serviço, então testem-no antes que os seus usuários o façam.
Faça streaming de bancos de dados operacionais com captura de dados de mudança
Muitas vezes vocês querem reagir a mudanças num banco de dados que nunca foi projetado para emitir eventos. A captura de dados de mudança (CDC) resolve isso lendo o log de transações do banco e transformando cada inserção, atualização e exclusão num fluxo de eventos de mudança. É muito melhor que consultar a tabela por um temporizador, o que é lento, perde estados intermediários e martela a origem. A CDC permite manter um índice de busca, um cache, um repositório analítico ou um serviço a jusante continuamente sincronizado com um sistema de registro, e faz isso sem mudanças invasivas na aplicação. Tratem o fluxo de mudanças como um produto de dados de primeira classe: versionem seu esquema, documentem seu significado e vigiem sua defasagem, porque tudo a jusante herda essa defasagem.
Prefira uma arquitetura streaming-first a manter duas bases de código
A clássica arquitetura Lambda roda uma camada batch para um histórico exato e completo ao lado de uma camada de velocidade para resultados frescos e aproximados e depois os mescla. Funciona, mas obriga vocês a escrever e manter a mesma lógica de negócio duas vezes, em dois sistemas, e a reconciliar as diferenças para sempre. A arquitetura Kappa colapsa isso: mantenham um log durável e reproduzível de eventos e rodem todo o processamento como processamento de fluxo, reprocessando o histórico ao reproduzir o log quando a lógica muda. O setor derivou para essa forma streaming-first porque uma única base de código é dramaticamente mais barata de manter e de raciocinar. Se vocês conseguem expressar as necessidades de batch como reproduções sobre um log de eventos retido, evitam por completo o imposto das duas bases de código. Usem brokers baseados em log que retêm o histórico para que o reprocessamento seja questão de rebobinar, não de reconstruir.
Exponha os fluxos como SQL, views materializadas e OLAP em tempo real
Nem todo mundo que precisa de streaming deveria ter que escrever código de processamento de fluxo de baixo nível. O SQL de streaming permite que analistas e engenheiros expressem janelas, junções e agregações numa linguagem que já conhecem e mantém os resultados continuamente atualizados como views materializadas. Para consultas analíticas de baixa latência sobre dados frescos, um repositório de processamento analítico online (OLAP) em tempo real ingere o fluxo e responde a consultas de fatiar e cortar em milissegundos, que é o que alimenta um painel operacional genuinamente ao vivo. Combinem isso com as práticas de análise de produto do capítulo 7.4 (análises de produto e experimentação) quando o objetivo é um retorno rápido sobre funcionalidades e experimentos. Escolham essas ferramentas de nível mais alto onde se ajustarem e guardem os processadores de fluxo escritos à mão para a lógica que elas não conseguem expressar.
Planeje a contrapressão e o reprocessamento desde o início
Um fluxo pode chegar mais depressa do que vocês conseguem processar. A contrapressão é o mecanismo que permite a um consumidor lento sinalizar a montante para desacelerar em vez de cair ou descartar dados em silêncio. Garantam que cada estágio do pipeline a respeite e monitorem a defasagem do consumidor como uma métrica de destaque, porque a defasagem crescente é o aviso mais precoce de que vocês estão perdendo a corrida. O reprocessamento é a outra capacidade que as pessoas gostariam de ter embutido. Quando vocês acham um bug ou mudam uma regra, querem reproduzir o histórico pela lógica corrigida. Isso só é possível se o log de eventos retém histórico suficiente e os destinos são idempotentes o bastante para absorver a reprodução. Projetem ambos desde o primeiro dia: adaptá-los sob pressão de incidente é miserável.
Compromissos: prós e contras
| Escolha | Prós | Contras | Melhor ajuste |
|---|---|---|---|
| Batch | Simples, barato, fácil de testar e reprocessar | Alta latência, obsoleto entre execuções | Relatórios, a maioria das análises |
| Micro-lote (minutos) | Quase em tempo real, muito mais simples que streaming | Não é verdadeiramente instantâneo | Painéis de “tempo real” |
| Streaming verdadeiro (subsegundo) | Reação instantânea, resultados contínuos | Complexo, caro, difícil de testar | Fraude, alertas, personalização ao vivo |
| Pelo-menos-uma-vez + destino idempotente | Resultados corretos, acessível, resiliente | Exige projeto disciplinado de chaves | A maioria dos pipelines de streaming |
| Maquinaria de exatamente-uma-vez | Garantia forte de ponta a ponta | Cara, limitada entre sistemas | Caminhos estreitos de alto risco |
| Lambda (batch + velocidade) | Histórico exato mais visão fresca | Duas bases de código a manter | Migrações de legado |
| Kappa (streaming-first) | Uma base de código, reproduzível | Exige log retido e durável | Novas plataformas de streaming |
A tensão central é latência contra complexidade. Cada passo rumo ao tempo real custa em carga operacional, dificuldade de teste e dinheiro, e os retornos não são lineares: passar de diário para a cada poucos minutos é barato e muitas vezes suficiente, enquanto passar de minutos para subsegundo é onde a despesa se concentra. Resolvam a tensão precificando a decisão, não a tecnologia. Perguntem que ação a frescura viabiliza e quanto custa o atraso e comprem só a redução de latência que essa ação justifica. Quando precisarem de streaming, apoiem-se na entrega pelo-menos-uma-vez com destinos idempotentes e um log streaming-first, porque essa combinação dá correção e reprodutibilidade sem as garantias mais pesadas.
Perguntas para discutir com sua equipe
Que decisão os dados em tempo real de fato viabilizam para nós, e quanto custa quando esses dados chegam um minuto atrasados em vez de instantaneamente? Esta é a pergunta que deve condicionar todo projeto de streaming, porque o streaming aproximadamente dobra o custo e a complexidade operacional em comparação com o batch. Uma grande equipe pode queimar trimestres construindo uma plataforma em tempo real que serve painéis que um humano confere duas vezes por dia, o que é dinheiro posto no fogo. Levem a ação concreta que os dados conduzem, seja bloquear uma transação fraudulenta, acionar um operador ou mudar o que um usuário vê, e ponham um número no custo da latência de cada uma. Se a resposta honesta é que um micro-lote de cinco minutos serviria à necessidade, isso é uma descoberta a celebrar, não a esconder. A resposta deve mudar diretamente se vocês constroem streaming verdadeiro, se contentam com micro-lotes ou ficam no batch.
Como tratamos os eventos atrasados e fora de ordem, e o que acontece com os dados que chegam depois de uma janela fechar? Os dados atrasados e fora de ordem são a parte difícil do streaming, e as equipes que pulam esta pergunta a descobrem em produção, quando os números se recusam a reconciliar. As pressões concorrentes são latência e correção: manter as janelas abertas por mais tempo para pegar retardatários atrasa todo resultado e consome mais memória, fechá-las mais depressa descarta dados reais em silêncio. Levem evidências sobre quão tarde os seus dados de fato chegam, medidas como a diferença entre o tempo do evento e o tempo de processamento em suas fontes, já que uma fonte móvel em túneis se comporta de modo muito diferente de um evento do lado do servidor. Decidam explicitamente se os dados atrasados são descartados, registrados ou disparam uma correção e garantam que todos a jusante saibam qual. Num contexto governamental em que os números precisam ser defensáveis, descartar eventos atrasados em silêncio pode ser um problema de conformidade, então a política precisa ser deliberada e documentada.
Os nossos destinos são idempotentes o bastante para reproduzirmos o histórico com segurança, e o nosso log de eventos retém o suficiente para tornar a reprodução possível? O reprocessamento é a capacidade que as equipes mais frequentemente gostariam de ter embutido e mais frequentemente não embutiram, e depende de duas coisas funcionando juntas: destinos idempotentes que absorvem eventos reproduzidos sem duplicar e um log durável que retém histórico suficiente para reproduzir. Sem ambos, consertar um bug de lógica significa que vocês não conseguem recalcular limpamente o período afetado e ficam presos remendando números à mão sob pressão. Levem a sua janela atual de retenção e um teste concreto: escolham um bug real do último trimestre e perguntem se teriam conseguido reproduzir a lógica corrigida sobre os dados afetados. A atração contrária é o custo, já que reter histórico e projetar escritas idempotentes exige armazenamento e disciplina de antemão. Mas a alternativa surge no pior momento possível, durante um incidente, então a resposta molda quanto vocês investem na reprodutibilidade antes de precisar dela.
Quando um job de streaming cai, com que rapidez ele deve se recuperar, quanto estado ele pode guardar e nós já cronometramos de fato uma recuperação sob carga de produção? Um job batch que morre pode ser reexecutado amanhã, mas um fluxo sempre ligado que morre é uma interrupção em andamento, e os jobs com estado que guardam contagens correntes, junções ou modelos de fraude podem perder minutos de memória ou levar muito tempo para recarregar o estado após uma reinicialização. Para uma grande equipe, é aqui que um detalhe sem glamour define em silêncio a sua disponibilidade real: o estado ilimitado cresce até o job esgotar a memória, e uma restauração lenta de checkpoint transforma um soluço de dez segundos num de dez minutos. As pressões concorrentes são frescor contra segurança, porque checkpoints mais frequentes encurtam a recuperação mas acrescentam sobrecarga, e uma retenção generosa de estado melhora a exatidão mas arrisca esgotar a memória. Levem um objetivo concreto de tempo de recuperação, o tamanho atual do estado e sua curva de crescimento, o intervalo dos seus checkpoints e os resultados de um exercício real de failover e não uma estimativa esperançosa. Em contextos corporativos e governamentais em que o fluxo sustenta a pontuação de fraude ou um canal de segurança pública, um caminho de recuperação não testado é um risco operacional aceito sem medir, então tratem o exercício como uma exigência, não um luxo.
Rodamos uma base de código streaming-first ou camadas separadas de batch e de velocidade, e quanto custa de fato mantê-las reconciliadas? O padrão Lambda, com uma camada batch para o histórico exato mais uma camada de velocidade para resultados frescos, obriga a escrever a mesma lógica de negócio duas vezes, em dois sistemas, e reconciliar as respostas para sempre, ao passo que uma forma streaming-first (Kappa) mantém um log durável e reproduzível e roda todo o processamento como processamento de fluxo. Para uma grande organização, a lógica duplicada é onde nascem a deriva e os números contestados, porque uma regra muda numa camada e não na outra, e os engenheiros gastam tempo real explicando por que as duas discordam. A atração para manter ambas é a inércia e o conforto de uma camada batch comprovada, então pesem isso contra o imposto de manutenção com honestidade. Levem a lista das computações que hoje rodam nos dois lugares, os incidentes causados pelas duas camadas discordando e uma avaliação de se o log de eventos retém histórico suficiente para expressar as necessidades de batch como reproduções. Em contextos governamentais e corporativos auditados, duas camadas que podem relatar números diferentes para o mesmo período são em si um passivo de conformidade, já que vocês precisam poder dizer qual número é autoritativo e por quê.
Quem opera este sistema sempre ligado quando ele quebra às três da manhã, e vocês orçaram a carga de sobreaviso e as habilidades especializadas que ele exige, ou estão presumindo uma equipe com formato de batch? O streaming desloca o custo da construção para a operação: o sistema precisa se manter saudável a cada segundo, o que significa cobertura real de sobreaviso, engenheiros fluentes em tempo de evento, marcas d’água, estado e semântica de entrega, e testes mais difíceis que os de um job que roda e para. As equipes rotineiramente aprovam uma plataforma de streaming pelas suas capacidades e nunca financiam as pessoas que a mantêm viva, então a plataforma se degrada e a confiança se erode. O compromisso é escopo contra sustentabilidade: cada pipeline adicional em tempo real é mais uma coisa que pode acionar alguém, então a pergunta é se a latência que ele compra justifica um compromisso operacional permanente. Levem um inventário honesto de quem é dono de cada fluxo em produção, o seu rodízio atual de sobreaviso e sua folga e onde de fato está a expertise em tempo de evento, seja uma contratação, um parceiro ou um serviço gerenciado. Para um órgão público ou uma grande empresa, acrescentem os prazos de contratação e de recrutamento e qualquer opção de serviço gerenciado, porque uma plataforma em tempo real que depende de talento escasso que vocês não conseguem recrutar nem reter é um plano de operar um sistema propenso a interrupções com equipe insuficiente.
Perspectiva por setor
Startup. O streaming raramente é o primeiro movimento, e montar uma plataforma pesada pode afundar uma equipe minúscula. Escolham o único sinal que toca o seu valor central, ponham os eventos num único broker baseado em log com retenção e rodem um processador leve com destinos com chave e idempotentes para que uma nova tentativa pelo-menos-uma-vez nunca conte em dobro. Mantenham alguns dias de histórico para poder reproduzir pela lógica corrigida e prefiram um serviço de streaming gerenciado a operar o próprio cluster, porque o seu recurso mais escasso é a atenção da engenharia.
Pequena empresa. Vocês provavelmente não têm especialista em streaming nem apetite para rodar infraestrutura sempre ligada, então tratem o tempo real como algo que se compra dentro de ferramentas que já usam e não como um sistema que se mantém com equipe. Enquadrem a necessidade como uma pergunta de latência com um número anexado, e na maioria dos casos um micro-lote que atualiza a cada poucos minutos a atenderá por uma fração do custo e do risco. Escolham fornecedores cujas funcionalidades em tempo real sejam transparentes sobre a defasagem e fáceis de abandonar e reservem o streaming sob medida para o raro caso em que dados frescos conduzem diretamente receita ou segurança.
Grande empresa. O problema é consistência e custo entre muitas equipes: uma plataforma compartilhada baseada em log, uma política padrão de tempo de evento e de dados atrasados e destinos idempotentes para que os grupos parem de reinventar pipelines frágeis. Orcem explicitamente a operação sempre ligada e a carga de sobreaviso, padronizem num log streaming-first para evitar uma base de código batch duplicada e gerenciem os fluxos como produtos de dados governados, com donos, versionamento de esquema e defasagem monitorada e não como um amontoado de jobs sob medida. Acompanhem a latência, o tempo de recuperação e o custo por fluxo como métricas de portfólio.
Governo. A auditabilidade e a prestação de contas públicas moldam toda escolha. Retenham cada evento processado num log durável para que os números reportados aos órgãos de supervisão, ocupação do transporte, anomalias de benefícios, decisões de fraude, possam ser reconstruídos exatamente, e tornem explícita e documentada a política de dados atrasados em vez de descartar eventos em silêncio. A contratação deve exigir portabilidade de dados e divulgação das garantias de entrega e de retenção de um serviço gerenciado, e qualquer retificação após uma mudança de regra deve ser uma reprodução defensável pela lógica corrigida, não um remendo manual que ninguém consegue rastrear.
Exemplos
Startup. Um aplicativo de consumo quer mostrar aos usuários um feed de atividade ao vivo e sinalizar logins suspeitos assim que acontecem. A equipe resiste a montar uma plataforma pesada de streaming. Ela põe os eventos num único broker baseado em log com retenção, roda um processador de fluxo leve para a lógica de risco de login e alimenta um repositório OLAP em tempo real que sustenta o feed de atividade. Todo destino tem chave e é idempotente, de modo que uma nova tentativa pelo-menos-uma-vez nunca conta em dobro. Quando mais tarde acha um bug na regra de risco, simplesmente reproduz o log pela lógica corrigida durante a noite, porque guardou uma semana de histórico e nunca precisou de uma segunda base de código batch.
Grande empresa. Um banco de varejo pontua cada transação de cartão quanto à fraude dentro da janela de autorização, juntando o fluxo ao vivo de transações a um modelo com estado do comportamento recente da conta. O checkpoint permite ao serviço de pontuação se recuperar de uma falha de nó em segundos sem perder a memória dos últimos minutos. Separadamente, a captura de dados de mudança leva atualizações do banco de dados bancário central para um índice de busca e um serviço de personalização, mantendo ambos frescos sem consultas periódicas. Os painéis operacionais leem de um repositório OLAP em tempo real para que as equipes de risco e de operações acompanhem o negócio em movimento, e todo o pipeline emite a telemetria de defasagem e de vazão descrita no capítulo 9.2.
Governo. Uma autoridade metropolitana de transporte ingere posições de veículos e toques de bilhete para prever chegadas e monitorar a lotação em tempo real, alimentando tanto aplicativos públicos quanto um centro de operações. Como os passageiros em túneis enviam os toques em rajadas atrasadas, a equipe calcula a ocupação pelo tempo do evento, com marcas d’água ajustadas ao atraso observado, e registra qualquer evento que chegue depois de sua janela fechar em vez de descartá-lo em silêncio. Cada evento processado é retido num log auditável para que os números de ocupação reportados aos órgãos de supervisão possam ser reconstruídos exatamente. Quando uma regra de tarifa muda, ela reproduz o período afetado pela lógica corrigida e produz uma retificação defensável.
Justificativa de negócio: motivações, ROI e TCO
O retorno dos dados em tempo real vem de agir enquanto a ação ainda importa. A fraude pega durante a autorização evita uma perda que um batch noturno apenas relataria. A personalização que responde dentro de uma sessão eleva a conversão de um modo que a recomendação de amanhã não consegue. O monitoramento operacional que reflete o presente permite intervir antes que um pequeno problema vire uma interrupção ou um incidente público. Em cada caso, o valor é o delta entre agir agora e agir depois, e esse delta é o que vocês devem quantificar ao defender o caso.
O custo total de propriedade é maior que o do batch, e a honestidade sobre isso protege a sua credibilidade. Vocês pagam por infraestrutura sempre ligada, por engenheiros que entendem tempo de evento, marcas d’água, estado e semântica de entrega e pela carga mais difícil de teste e de sobreaviso de um sistema que precisa se manter saudável continuamente em vez de rodar e parar. Uma arquitetura streaming-first sobre um log retido reduz o custo contínuo ao poupar vocês de uma base de código batch duplicada, e escolher pelo-menos-uma-vez com destinos idempotentes evita a despesa da maquinaria de exatamente-uma-vez de ponta a ponta. O erro mais caro é construir tempo real onde o micro-lote ou o batch bastaria, então o argumento de custo mais forte é muitas vezes uma decisão de não fazer streaming. Enquadrem o argumento à liderança em torno de decisões específicas sensíveis à latência e do seu retorno mensurável e sejam igualmente claros sobre onde ficar no batch economiza dinheiro sem perda de valor.
Antipadrões e armadilhas
- Construir streaming por prestígio quando um micro-lote a cada poucos minutos atenderia à necessidade.
- Calcular sobre o tempo de processamento, de modo que os números oscilam com a infraestrutura em vez de com o mundo.
- Ignorar os dados atrasados e fora de ordem até a reconciliação falhar em produção.
- Perseguir o exatamente-uma-vez literal em toda parte em vez de pelo-menos-uma-vez com destinos idempotentes.
- Estado ilimitado sem expiração, crescendo em silêncio até um job esgotar a memória.
- Nenhum checkpoint, de modo que uma reinicialização perde o estado ou força uma reprodução completa.
- Consultar bancos de dados operacionais por temporizador em vez de usar a captura de dados de mudança.
- Manter uma camada batch e uma camada de velocidade Lambda com lógica duplicada e à deriva.
- Uma janela de retenção curta demais para reproduzir o histórico quando se acha um bug.
- Fluxos sem métricas de defasagem, vazão ou frescor, falhando em silêncio.
Modelo de maturidade
- Nível 1, Iniciar: Tudo é batch, ou alguns jobs de streaming feitos à mão rodam de modo reativo, sem monitoramento. Os números são calculados sobre o tempo de processamento, os dados atrasados são ignorados e uma reinicialização perde o estado. Ninguém consegue reproduzir o histórico para consertar um bug, e os problemas são descobertos quando os números a jusante se recusam a reconciliar.
- Nível 2, Desenvolver: Algumas equipes rodam pipelines centrais de streaming num broker baseado em log, com checkpoints, e distinguem o tempo do evento do tempo de processamento e usam janelas básicas. A prática é inconsistente de equipe para equipe: a entrega é pelo-menos-uma-vez mas nem todos os destinos são idempotentes, o tratamento de dados atrasados é improvisado e a defasagem é vigiada informalmente em vez de alertada.
- Nível 3, Padronizar: O tempo do evento, as marcas d’água e uma política explícita de dados atrasados são documentados e aplicados em toda a organização. Os destinos são idempotentes para resultados efetivamente-uma-vez, o estado tem expiração e a captura de dados de mudança alimenta os sistemas a jusante por convenção. Um log retido sustenta a reprodução, e a defasagem, a vazão e o frescor são monitorados com alertas como padrão da organização e não como hábito de cada equipe.
- Nível 4, Gerenciar: O patrimônio de streaming é medido e controlado em relação a linhas de base. Cada pipeline carrega objetivos de nível de serviço para latência de ponta a ponta, defasagem do consumidor, tempo de recuperação, distorção do tempo de evento, taxa de eventos atrasados, tamanho do estado e custo por milhão de eventos, todos acompanhados contra metas combinadas e com alertas de regressão. A recuperação é exercitada e cronometrada e não presumida, a folga de contrapressão e o crescimento do estado são vigiados como sinais de capacidade, e um novo fluxo precisa cruzar essas métricas antes de ir para produção.
- Nível 5, Orquestrar: Uma arquitetura streaming-first atende tanto às necessidades frescas quanto às históricas a partir de um log reproduzível, e o SQL de streaming, as views materializadas e o OLAP em tempo real tornam os dados frescos amplamente acessíveis. O reprocessamento é rotineiro e testado, a plataforma escala automaticamente e se rebalanceia conforme a carga e o custo medidos, e os fluxos são aposentados, redefinidos ou substituídos com base em evidências. O streaming é integrado ao planejamento de negócio e de risco, e todo fluxo é observável e auditável de ponta a ponta à medida que o quadro de carga e de custo muda.
Ideias para discussão
- Onde, na sua pilha, o “tempo real” realmente merece o seu custo, e onde é um desejo não examinado?
- Quão grande é a diferença entre o tempo do evento e o tempo de processamento em suas fontes, e vocês a medem?
- Vocês conseguiriam colapsar uma configuração Lambda de batch e velocidade numa única base de código streaming-first, e o que impediria?
- Quais dos seus destinos são de fato idempotentes, e vocês conseguiriam reproduzir hoje com segurança os dados do último trimestre pela lógica corrigida?
- Qual é a sua política para os dados que chegam depois de uma janela fechar, e todos a jusante a conhecem?
- Como a captura de dados de mudança alteraria o modo como vocês mantêm busca, caches e análises sincronizados?
Principais conclusões
- Recorram ao streaming apenas quando uma decisão sensível à latência o paga. O batch e o micro-lote são padrões mais baratos.
- Calculem sobre o tempo do evento e tratem os dados atrasados e fora de ordem como o problema central, tratado com janelas e marcas d’água.
- Prefiram a entrega pelo-menos-uma-vez com destinos idempotentes para resultados efetivamente-uma-vez a um exatamente-uma-vez literal em toda parte.
- Façam checkpoints do processamento com estado, limitem o estado e monitorem a defasagem do consumidor como uma métrica de destaque.
- Usem a captura de dados de mudança para fazer streaming de bancos de dados operacionais em vez de consultá-los por temporizador.
- Favoreçam uma arquitetura streaming-first sobre um log retido e reproduzível a manter duas bases de código.
- Exponham os fluxos por SQL de streaming, views materializadas e OLAP em tempo real e mantenham todo fluxo observável e auditável.
Referências e leitura complementar
- Tyler Akidau, Slava Chernyak, and Reuven Lax, “Streaming Systems.”
- Martin Kleppmann, “Designing Data-Intensive Applications.”
- Nathan Marz and James Warren, “Big Data” (Lambda architecture).
- Jay Kreps, “Questioning the Lambda Architecture” (O’Reilly Radar).
- Fabian Hueske and Vasiliki Kalavri, “Stream Processing with Apache Flink.”
- Ben Stopford, “Designing Event-Driven Systems.”
- Tyler Akidau and colleagues, “The Dataflow Model” (VLDB paper on windowing and watermarks).