Google PDE: Mensageria, Ingestão de Eventos e Serviços em Tempo Real — 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
Os serviços de mensageria, ingestão de eventos e tempo real no Google Cloud se concentram no Cloud Pub/Sub e no Eventarc para transporte desacoplado e durável; no Dataflow para processamento de stream com estado (stateful); e em sinks como BigQuery, Cloud Storage e bancos de dados operacionais. Projetar para entrega do tipo “pelo menos uma vez” (at-least-once), consumo idempotente e observabilidade garante sistemas resilientes que escalam elasticamente, mantendo a correção sob falhas, contrapressão (backpressure) e evolução de esquema.
Mensageria Principal com o Pub/Sub
- Tópicos e inscrições (subscriptions)
- Publishers enviam mensagens para um tópico; subscribers se anexam por meio de inscrições (vários subscribers podem consumir as mesmas mensagens de forma independente).
- Tipos de inscrição:
- Pull: os clientes buscam explicitamente as mensagens; use o streaming pull para obter o maior throughput e menos viagens de ida e volta (round trips).
- Push: o Pub/Sub entrega via HTTPS; seu endpoint deve retornar 2xx para confirmar o recebimento (acknowledge).
- Exportar para o BigQuery: uma inscrição do tipo BigQuery entrega mensagens para uma tabela do BigQuery sem código; é a melhor opção quando os payloads correspondem ao esquema declarado e a ingestão de baixa latência em analytics é necessária.
- Chaves de ordenação (ordering keys)
- Habilite a ordenação de mensagens no tópico e na inscrição para receber entregas em ordem por chave de ordenação. O throughput por chave é serializado: uma mensagem em trânsito (in-flight) por chave pode bloquear as subsequentes; use muitas chaves (por exemplo, hash(device_id)) para escalar.
- Fan-out e replay
- Crie inscrições separadas para diferentes consumidores para isolar as cargas de trabalho (workloads) e a retenção.
- Use “seek” ou “snapshot” para fazer o replay a partir de um timestamp ou snapshot para recuperação e preenchimentos retroativos (backfills).
Trade-offs:
- A ordenação reduz o paralelismo e o throughput por chave; desabilite a ordenação a menos que seja estritamente necessário.
- O modo Push simplifica o código do cliente, mas introduz preocupações com o escalonamento do endpoint HTTP, segurança e backoff; o modo Pull oferece mais controle e estabilidade em alto throughput.
Semântica de Entrega, Confirmação (Acknowledgment), Retenção e Dead Lettering
- Confirmação (acknowledgment) e prazos (deadlines)
- Entrega do tipo “pelo menos uma vez” (at-least-once): duplicatas podem ocorrer.
- Cada entrega tem um prazo de confirmação (ack deadline) (padrão de 10 segundos). Estenda-o (ModifyAckDeadline) ao processar trabalhos de longa duração; a falha em confirmar (ack) antes do prazo é a causa mais comum de entregas push duplicadas.
- Um Nack (confirmação negativa) ou a expiração do prazo torna a mensagem elegível para reentrega.
- Retenção
- Mensagens não confirmadas (unacknowledged) são retidas pelo prazo de confirmação da inscrição e reenviadas; mensagens confirmadas podem ser retidas por até a duração de retenção de mensagens do tópico para replay. Configure a retenção para cobrir sua interrupção máxima mais o tempo de recuperação.
- Novas tentativas (retries)
- Pull: a reentrega ocorre após a expiração do ack deadline; controle a concorrência com limites de controle de fluxo (flow control).
- Push: backoff exponencial; apenas HTTP 2xx é sucesso. 3xx/4xx/5xx acionam novas tentativas. Implemente handlers idempotentes para tolerar repetições.
- Tópicos de mensagens mortas (Dead-letter topics - DLTs)
- Configure um tópico DL e um número máximo de tentativas de entrega por inscrição para colocar mensagens “venenosas” (poison messages) em quarentena.
- Monitore o volume da DLQ (Dead Letter Queue); crie fluxos de trabalho de triagem e republique no tópico principal após a correção.
Exemplo:
undefined
Resumo da semântica de entrega:
- Pub/Sub: entrega do tipo “pelo menos uma vez” (at-least-once), ordenação de melhor esforço (best-effort) dentro de uma chave de ordenação, se habilitada.
- Sinks: as APIs de inserção do BigQuery fornecem mitigação de duplicatas (insertId ou offsets de stream da Storage Write API), mas ainda assim projete consumidores e gravadores para serem idempotentes.
Esquemas (Schemas), Compatibilidade e Validação
- Esquemas do Pub/Sub
- Suporte nativo para Avro e Protocol Buffers com esquemas armazenados centralmente.
- Configurações de esquema no nível do tópico: codificação (Avro ou Protobuf) e imposição (enforcement) (nenhuma, apenas validar ou exigir).
- O produtor publica payloads codificados; o Pub/Sub valida em relação ao esquema atual quando a imposição está habilitada.
- Evolução e compatibilidade
- Use alterações retrocompatíveis (backward-compatible) (adicionar campos opcionais, adicionar campos com valores padrão em Avro, nunca reutilizar tags em Protobuf, evitar remover ou renomear campos).
- Versione os esquemas explicitamente. Para alterações que quebram a compatibilidade (breaking changes), faça publicação dupla (dual-publish) para tópicos v1 e v2, ou adicione um campo de versão e roteie de acordo.
- Contratos produtor-consumidor
- Os consumidores devem ignorar campos desconhecidos e usar valores padrão para os ausentes.
- Teste a compatibilidade do esquema em todos os consumidores antes da promoção para produção; valide em inscrições de homologação (staging) com a mesma imposição de esquema da produção.
Exemplo curto de Avro (trecho):
undefined
Padrões de Ingestão de Streaming, Vazão, Escalabilidade, Segurança e Operações
- Padrões de ingestão em tempo real
- Pub/Sub -> Dataflow -> BigQuery: use o sink da BigQuery Storage Write API para alta vazão e idempotência com offsets de stream; direcione falhas para uma tabela de dead-letter para inspeção.
- Pub/Sub -> Dataflow -> Cloud Storage: arquive eventos brutos para reprocessamento; use escritas em janela e compactadas para equilibrar custo e latência.
- Pub/Sub -> armazenamentos operacionais: escreva no Bigtable para buscas de baixa latência, no Spanner para transações fortemente consistentes, ou no Cloud SQL/Firestore com base nas necessidades da carga de trabalho. Garanta upserts idempotentes com chave por um ID de evento único.
- Entrega pelo menos uma vez (at-least-once), prevenção de duplicatas e idempotência
- Transporte um
event_ideevent_timeúnicos em cada mensagem; imponha UUIDs no lado do produtor. - Deduplicação de streaming do BigQuery: defina o
insertIdou use a Storage Write API com streams ordenados; ainda assim, proteja as consultas com lógica de deduplicação. - Exemplo de deduplicação em tempo de consulta: WITH ranked AS ( SELECT t.*, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_time DESC) AS rn FROM dataset.events t ) SELECT * EXCEPT(rn) FROM ranked WHERE rn = 1;
- Para endpoints de push, retorne 2xx somente após o processamento bem-sucedido; caso contrário, espere uma reentrega.
- Transporte um
- Vazão de mensagens, cotas e escalabilidade
- Publishers: envie mensagens em lote e reutilize conexões; paralelize entre múltiplos clientes. Use muitas chaves de ordenação (ordering keys) para escalar cargas de trabalho ordenadas.
- Subscribers: prefira streaming pull com controle de fluxo (máximo de bytes/mensagens pendentes). Dimensione os prazos de confirmação (ack deadlines) de acordo com o tempo de processamento e estenda-os quando necessário.
- Monitore e solicite aumentos de cota para a vazão de publicação e assinatura conforme os volumes crescem; projete com uma margem de segurança (por exemplo, 2x o pico esperado) para absorver picos de tráfego.
- Consistência e disponibilidade
- O streaming do BigQuery é eventualmente consistente para a visibilidade em consultas; para consultas interativas que devem incluir linhas transmitidas, aguarde com base na latência observada (por exemplo, 2x o atraso de disponibilidade P50) ou projete usando agregações alinhadas a watermarks no Dataflow e consulte os resultados materializados.
- Segurança
- IAM: conceda papéis de menor privilégio (
pubsub.publisherpara produtores no tópico;pubsub.subscriberpara consumidores na assinatura). Use contas de serviço (service accounts) dedicadas para cada carga de trabalho. - Autenticação de push: configure assinaturas de push para anexar tokens OIDC de uma conta de serviço; imponha a validação de público (audience) no endpoint. Prefira endpoints privados do Cloud Run para autenticação e TLS integrados.
- Criptografia: o Pub/Sub criptografa em trânsito e em repouso; use CMEK em tópicos para chaves gerenciadas pelo cliente. Aplique VPC Service Controls para reduzir o risco de exfiltração de dados. Use criptografia do lado do cliente para campos de payload sensíveis, se necessário.
- IAM: conceda papéis de menor privilégio (
- Diagnóstico operacional de lag, reentrega e falha de subscriber
- Monitore com o Cloud Monitoring:
subscription/num_undelivered_messageseoldest_unacked_message_agepara o backlog.expired_ack_deadline_countpara detectar acks perdidos que causam duplicatas.publish_request_countepull_request_countpara a vazão.
- Investigue eventos ausentes no dashboard ao reproduzir um conjunto de dados conhecido através do pipeline e comparar as saídas de cada estágio para isolar a transformação ou o sink com falha.
- Para streaming com Dataflow:
- Use o autoescalonamento com um
maxWorkersapropriado para absorver a carga de muitas fontes. - Drene os pipelines (
drain) para atualizações incompatíveis para permitir que o trabalho em andamento seja concluído e evitar a perda de dados.
- Use o autoescalonamento com um
- Para notificações de inserção do BigQuery, direcione as entradas de auditoria do Cloud Logging através de um sink filtrado para tabelas específicas para um tópico do Pub/Sub para alertas.
- Monitore com o Cloud Monitoring:
Cenário de Problema Prático
A Contoso Freight precisa de uma plataforma global de eventos em tempo real para ingerir 10.000 mensagens de telemetria de IoT por minuto de caminhões, enriquecer eventos, potencializar análises interativas e acionar fluxos de trabalho na chegada de arquivos de parceiros externos. Alguns CSVs de parceiros contêm linhas malformadas, e a equipe de análise precisa inspecionar os erros sem bloquear o stream.
- Criar a camada principal de mensagens e esquema
- Ação: Defina um esquema Avro para telemetria e anexe-o a um tópico do Pub/Sub
telemetrycom a validação de esquema (schema enforcement) definida comorequire. Habilite a ordenação de mensagens e publique comordering_key = hash(device_id). - Justificativa: A validação de esquema no nível do tópico rejeita eventos malformados precocemente. A ordenação por dispositivo suporta o processamento ordenado quando necessário, enquanto o hashing distribui as chaves para preservar a vazão.
- Provisionar assinaturas com isolamento e dead-lettering
- Ação: Crie uma assinatura pull
telemetry-stream-subpara o Dataflow com um tópico de dead-lettertelemetry-dltemax_delivery_attempts=10. Adicione uma assinatura do BigQuerytelemetry-raw-bqpara armazenar eventos brutos em uma tabela particionada por tempo para linhagem e repetição. - Justificativa: A DLQ (Dead-Letter Queue) isola as poison messages para investigação. Uma assinatura separada do BigQuery fornece um caminho de exportação de baixa sobrecarga operacional para a preservação de eventos brutos, independente do pipeline de processamento.
- Construir um pipeline de streaming do Dataflow para enriquecimento e sinks
- Ação: Faça a ingestão a partir de
telemetry-stream-subusando streaming pull com controle de fluxo. Valide em relação ao esquema, enriqueça com dados de referência e calcule agregados em janela. Escreva no BigQuery usando a Storage Write API com um stream nomeado einsertId = event_id; escreva backups brutos no Cloud Storage a cada hora; redirecione registros ruins/com falha para uma tabela de dead-letter no BigQuery. - Justificativa: A Storage Write API proporciona escritas de alta vazão e baixa latência com idempotência via
insertId/offsets de stream. Uma tabela de dead-letter permite a inspeção sem bloquear o stream, e os arquivos no Cloud Storage permitem a repetição.
- Lidar com duplicatas e consistência eventual em análises
- Ação: Para consultas interativas que devem excluir duplicatas, publique
event_ideevent_timeem cada registro e use uma view de deduplicação: CREATE OR REPLACE VIEW analytics.latest_events AS SELECT * EXCEPT(rn) FROM ( SELECT e.*, ROW_NUMBER() OVER (PARTITION BY event_id ORDER BY event_time DESC) rn FROM analytics.events e ) WHERE rn = 1; Introduza um breve atraso na consulta com base na disponibilidade de streaming observada do BigQuery (por exemplo, o dobro da latência mediana). - Justificativa: A entrega do tipo at-least-once requer escritas idempotentes e deduplicação em tempo de consulta. A espera reduz a omissão de linhas em trânsito, dada a latência de visibilidade do streaming.
- Integrar a chegada de arquivos de parceiros com o Eventarc
- Ação: Configure o Eventarc para rotear eventos
object.finalizeddo Cloud Storage para o bucketpartner-dropspara um serviço do Cloud Run que inicia um job em lote do Dataflow para carregar CSVs no BigQuery, enviando erros de parsing para uma tabela de dead-letter. - Justificativa: O Eventarc fornece orquestração orientada a eventos com filtragem de CloudEvents por bucket e prefixo de objeto. Um job em lote do Dataflow separa as linhas malformadas para análise enquanto carrega os dados bons prontamente.
- Proteger a plataforma
- Ação: Use contas de serviço (service accounts) distintas: os produtores recebem
pubsub.publisheremtelemetry; a SA do worker do Dataflow recebepubsub.subscriberemtelemetry-stream-sube acesso de escrita nos datasets do BigQuery e no Cloud Storage de destino; o gatilho do Eventarc usa uma SA dedicada com o papel deinvokerno Cloud Run. Habilite o CMEK no tópicotelemetrye nos datasets do BigQuery. Configure endpoints de push, se houver, com OIDC e verificações de público (audience). - Justificativa: O IAM de menor privilégio e o CMEK atendem aos requisitos de segurança e conformidade; a entrega autenticada impede a falsificação (spoofing).
- Operar e escalar com confiabilidade
- Ação: Configure o autoescalonamento do Dataflow com um
maxWorkersgeneroso para absorver picos. Monitoresubscription/oldest_unacked_message_ageeexpired_ack_deadline_count; alerte quando os limites forem excedidos. Para alterações de pipeline que quebram a compatibilidade, faça o deploy com a opção de drenagem (drain) para evitar a perda de mensagens. Se o lag aumentar, aumente o paralelismo do subscriber e estenda os prazos de confirmação (ack deadlines) proporcionalmente ao tempo de processamento. - Justificativa: O monitoramento proativo detecta lag e reentregas precocemente. O autoescalonamento e os prazos de confirmação (ack deadlines) ajustados evitam tempestades de duplicatas. A drenagem (draining) preserva as mensagens em trânsito durante as atualizações.
Este design oferece uma ingestão em tempo real resiliente, segura e observável com integração de lote orientada a eventos, suporta tolerância a duplicatas e evolução de esquema, e oferece análises rápidas enquanto isola dados ruins para remediação direcionada.
← Processamento de Streams com Dataflow e Apache Beam · Todos os domínios · Spark →
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 →