Google PDE: Processamento de Streams com Dataflow e Apache Beam — Guia de estudos

Faz parte do Google Professional Data Engineer — Guia de estudos. Pratique com respostas verificadas no centro de exames da Google, ou faça testes cronometrados no ExamRoll.io.

Visão Geral

O processamento de stream no Google Cloud se concentra no modelo de programação unificado do Apache Beam, executado pelo runner do Dataflow. O Beam fornece uma abstração lógica — pipelines de transformações sobre PCollections — que desacopla seu código de detalhes de execução como paralelismo, autoescalonamento e tolerância a falhas. Em streaming, a corretude depende da semântica de tempo (tempo do evento vs. tempo de processamento), janelamento (fixo, deslizante, de sessão, global), marcas d’água (watermarks), gatilhos (triggers) e tratamento de dados tardios. A excelência operacional no Dataflow exige o dimensionamento correto dos workers, política de autoescalonamento, streaming engine, escolhas de shuffle, design de sink idempotente, tratamento de dead-letter e observabilidade robusta.

Modelo Apache Beam e semântica de tempo

Modos de falha e trade-offs:

Implantação, templates e estratégias de upgrade

undefined

undefined

Cenário de Problema Prático

A NovaTrack Inc. ingere telemetria de IoT global de 50.000 sensores de temperatura e deve fornecer agregados em nível de minuto, persistir dados brutos e alimentar um dashboard em tempo real. Mensagens malformadas ocasionais e entrega fora de ordem são esperadas. A solução deve escalar automaticamente, expor registros inválidos para inspeção e suportar upgrades sem tempo de inatividade.

Abordagem:

  1. Ingestão e semântica de tempo

    • Crie um tópico regional do Pub/Sub e publicadores por região com os atributos deviceId e eventTs (RFC3339). Habilite chaves de ordenação por deviceId quando viável.
    • Justificativa: O Pub/Sub fornece uma entrada durável e elástica com entrega pelo menos uma vez (at-least-once). Anexar timestamps de evento na borda preserva o tempo real do evento; a ordenação por dispositivo reduz a reordenação intra-dispositivo sem gargalos centrais.
  2. Pipeline de streaming do Dataflow com janelas de tempo de evento

    • Leia de uma assinatura dedicada via PubSubIO, extraindo eventTs como o timestamp do Beam, recorrendo a publishTime se ausente.
    • Aplique FixedWindows de 1 minuto com um gatilho antecipado aos 30 segundos e disparos tardios para cada elemento atrasado; defina o atraso permitido (allowed lateness) para 10 minutos e painéis acumulativos (accumulating panes).
    • Justificativa: Janelas de tempo de evento garantem agregados de minuto precisos; disparos antecipados alimentam o dashboard com dados atualizados em menos de um minuto; disparos tardios corrigem os agregados à medida que dados atrasados chegam. O limite de atraso limita o tamanho do estado e o custo.
  3. Validação, enriquecimento e roteamento para dead-letter

    • Implemente um ParDo que analisa JSON, valida esquema e intervalos, e enriquece com pequenos dados de referência estáticos por meio de um side input carregado do BigQuery no início do job.
    • Use TupleTags para emitir registros válidos para a saída principal e falhas para uma PCollection de dead-letter contendo o payload, erro, deviceId e timestamp de análise; escreva a DLQ em uma tabela particionada do BigQuery.
    • Justificativa: Side inputs mantêm os dados de referência em memória para baixa latência. A captura de dead-letter permite a inspeção e o reprocessamento direcionado de linhas inválidas sem bloquear o fluxo principal.
  4. Agregação e mitigação de hot-key

    • Agrupe por deviceId e calcule a média/mín/máx por minuto com CombineFns. Para métricas regionais de top-N, fragmente por região#N para evitar hot keys e, em seguida, reagregue.
    • Justificativa: Combiners minimizam o volume e o custo do shuffle; o particionamento de chaves (key sharding) evita gargalos de chave única durante o fan-in regional.
  5. Sinks e efeitos exactly-once

    • Escreva eventos brutos validados e agregados de minuto no BigQuery usando BigQueryIO com a Storage Write API. Defina um insert id estável baseado em deviceId + eventTs para idempotência em quaisquer novas tentativas personalizadas.
    • Justificativa: A Storage Write API fornece ingestão de alta taxa de transferência e baixa latência com semântica exactly-once dentro de um stream. IDs estáveis garantem a deduplicação downstream se ocorrerem repetições.
  6. Estratégia de consistência do dashboard

    • O dashboard consulta tabelas de agregados particionadas com uma retrospectiva de 2 minutos em relação à marca d’água (watermark) ou um atraso fixo de 2x a latência de disponibilidade observada para dados de streaming.
    • Justificativa: A visibilidade de streaming do BigQuery é eventualmente consistente; adiar ligeiramente as leituras evita a perda de linhas em trânsito, mantendo um comportamento quase em tempo real.
  7. Operações: autoescalonamento e streaming engine

    • Habilite o Streaming Engine; defina maxWorkers com base no pico esperado (ex: 3x a média), selecione um tipo de máquina dimensionado para a análise e criptografia que consomem CPU, e aumente o disco de inicialização para acomodar o shuffle transiente.
    • Monitore o atraso da marca d’água (watermark lag), segundos de backlog, CPU e taxa de transferência por etapa; alerte sobre atrasos sustentados e picos na taxa de DLQ.
    • Justificativa: O Streaming Engine externaliza o estado/shuffle para elasticidade e upgrades mais simples; o dimensionamento correto e o monitoramento evitam violações silenciosas de SLO.
  8. Implantação e upgrades com Flex Templates

    • Empacote o pipeline como um Flex Template com parâmetros: assinatura de entrada, tabelas de saída, tabela de DLQ, maxWorkers e região. Para uma alteração incompatível, inicie o novo pipeline visando o mesmo tópico com uma nova assinatura, verifique as saídas e, em seguida, drene o job antigo. Opcionalmente, crie um snapshot do Pub/Sub e faça com que a nova assinatura busque a partir do snapshot para garantir que não haja lacunas.
    • Justificativa: Flex Templates permitem implantações repetíveis e parametrizadas. Uma transição blue/green verificada com drenagem alcança perda zero de dados e tempo de inatividade mínimo.
  9. Reprocessamento e backfills em lote

    • Armazene arquivos Avro compactados de eventos brutos no Cloud Storage por meio de uma saída secundária (side output); execute um pipeline de Dataflow em lote para fazer backfill ou reprocessar para o BigQuery quando modelos ou esquemas mudarem.
    • Justificativa: Arquivos brutos duráveis suportam a reprodutibilidade e a evolução do esquema sem impactar o caminho principal (hot path).

Este design produz agregados corretos e de baixa latência com custo limitado, isolamento claro de erros, forte observabilidade e caminhos de upgrade seguros, ao mesmo tempo em que lida com dados fora de ordem e atrasados em escala global.


Analytics com BigQuery e Engenharia de Data Warehouse · Todos os domínios · Mensageria

Pratique estas questões → · Prática cronometrada no ExamRoll.io →

Pass the whole exam — not just this question

You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.

Passe no seu exame →

Navegar Google →

Related guides

Acesso completo

Uma assinatura. Todos os exames.

Todo plano desbloqueia pesquisa ilimitada de respostas, testes práticos, explicações de AI e a biblioteca completa de recursos — em mais de 20 idiomas.

Mensal
24.87
Just €0.83/day
Tudo incluído:
  • Pesquisa ilimitada de respostas
  • Testes práticos ilimitados
  • Explicações com AI
  • Biblioteca completa de recursos
  • Mais de 20 idiomas
  • Atualizações semanais de conteúdo
  • Recompensas e indicações
  • Suporte prioritário
Iniciar teste grátis

*Não é necessário cartão de crédito

Melhor custo-benefício
12 meses
179.87
Just €0.49/daySave 40%
Tudo incluído:
  • Pesquisa ilimitada de respostas
  • Testes práticos ilimitados
  • Explicações com AI
  • Biblioteca completa de recursos
  • Mais de 20 idiomas
  • Atualizações semanais de conteúdo
  • Recompensas e indicações
  • Suporte prioritário
Iniciar teste grátis

*Não é necessário cartão de crédito

✓ Plano gratuito incluído · ✓ Cancele a qualquer momento · ✓ Todos os planos desbloqueiam o produto completo