← Professional Data Engineer: pipelines e decisões de dados
10 / 15 · 135 MIN

Ingestão, eventos tardios e reprocessamento

Define identidade e tempo dos eventos; pratica deduplicação, conflitos, fecho temporal e recuperação controlada.

1. Definir o contrato de ingestão

Antes de escolher serviços, descreve o que uma linha representa e quando pode ser usada. Um evento pode ser uma instrução imutável, uma versão de uma entidade ou uma correção a um movimento anterior. Esses modelos exigem regras diferentes. Num contrato de eventos imutáveis, repetir um identificador com o mesmo conteúdo é uma repetição; repetir esse identificador com outro montante é um conflito. Num contrato de alterações versionadas, a mesma entidade pode ter várias versões legítimas. Não elimines versões apenas porque partilham a chave da entidade. Regista identidade do evento, identidade da entidade, versão do esquema, tempo de negócio, tempo de publicação e unidade dos valores. Num exemplo APS fictício, uma equipa recebe movimentos de carteiras para preparar o fecho diário. A origem pode ficar indisponível e enviar o lote horas depois. O negócio precisa de saber se um movimento entrou no período correto, se foi contado duas vezes e se uma correção alterou um relatório já aceite. Prepara três resultados observáveis: registos aceites, registos rejeitados com motivo e reconciliação por período. Define quem decide sobre conflitos e correções. Uma fila sem backlog pode coexistir com movimentos incorretamente descartados. A aceitação deve medir a utilização correta dos dados, para além da velocidade da ingestão. Os exemplos desta aula são fictícios e não descrevem procedimentos BNP Paribas.

2. Tornar o arranque reproduzível

Separa identidade de submissão, identidade dos workers e acesso aos recursos. Uma leitura bem-sucedida com a tua conta não prova que os workers conseguem ler a mesma origem. Quando o erro identifica uma conta de serviço, confirma o principal e o recurso antes de aumentar permissões. Na submissão, verifica também se o principal autorizado pode associar a conta dedicada ao job. No diagnóstico, distingue recusa IAM, recurso inexistente, erro DNS e timeout de ligação. Pedir permissões adicionais para resolver um timeout sem investigar a rede pode aumentar o acesso e deixar a falha intacta. Mantém no plano a região dos workers, a sub-rede selecionada e a capacidade de endereços. Numa Shared VPC, a referência da sub-rede deve identificar o projeto host. Uma equipa que prevê crescer de 20 para 100 workers precisa de verificar os endereços livres com os responsáveis da rede. Regista as premissas e o resultado do ensaio de capacidade. Para a release, conserva a identidade da imagem, da especificação Flex Template e dos parâmetros aprovados. Uma tag que pode mudar não é evidência suficiente de qual executável passou os testes. O pacote de mudança deve permitir reconstruir o que foi submetido e relacioná-lo com os dados usados na validação.

3. Separar repetição, conflito e alteração legítima

A identidade de transporte e a identidade de negócio respondem a perguntas diferentes. O ID da mensagem ajuda a reconhecer a reentrega dessa mensagem. Se o produtor publica duas vezes a mesma instrução e recebe IDs diferentes, a identidade de negócio continua a ser necessária. Define o âmbito da chave: pode precisar de incluir sistema de origem, identificador da instrução e versão. Documenta quanto tempo a memória de deduplicação dura e o que acontece quando um replay ultrapassa esse prazo. Guardar apenas a última chave observada não resolve repetições intercaladas nem execução distribuída. Trata conteúdo incompatível sob a mesma identidade como evidência a conservar. Se A contém quantidade 7 e depois quantidade 9, descartar a segunda chegada como duplicado esconde uma diferença. Substituir automaticamente o primeiro conteúdo também seria uma regra nova. Encaminha o conflito com contexto suficiente para decidir, evitando expor dados sensíveis desnecessários. No destino, verifica o âmbito da proteção que usas. Os offsets do Storage Write API são posições de um stream; não são uma chave global de negócio entre streams diferentes. Numa resposta perdida, conserva stream, offset e conteúdo para uma tentativa coerente. A equipa deve distinguir a confirmação da escrita do estado funcional do movimento no relatório.

4. Decidir o que fazer com dados tardios

Tempo do evento é quando a ocorrência de negócio aconteceu; tempo de processamento é quando uma etapa a trata. Define qual alimenta o relatório. Um movimento de ontem publicado hoje pode continuar a pertencer ao fecho de ontem. Em Beam, as janelas agrupam eventos e os triggers determinam emissões. O watermark representa progresso estimado em tempo de evento, e a política de atraso condiciona a inclusão de dados que chegam depois. Exactly-once não prova completude perante eventos descartados por atraso. O contrato do consumidor precisa de distinguir resultado provisório, fecho técnico e aceitação de negócio. No nosso caso, o relatório publicado deve manter histórico. A equipa acorda uma tolerância para a operação normal e um processo separado para correções fora desse prazo. Guarda os eventos excluídos, o motivo, o período afetado e a decisão tomada. Decide como uma versão corretiva será identificada e como os consumidores serão avisados. Não aumentes a tolerância sem avaliar estado mantido, latência e custo. Também não interpretes uma janela fechada como autorização para eliminar a evidência. Um bom ensaio inclui chegada fora de ordem, atraso dentro e fora da tolerância, correção aprovada e repetição de um evento antigo. Compara o resultado esperado pelo negócio com o que cada política realmente inclui.

5. Preparar um replay sem repetir efeitos

Um plano de replay começa pela disponibilidade dos dados e pelo âmbito a recuperar. Verifica retenção, fronteira inicial, seleção e destino. Um snapshot Pub/Sub guarda estado de confirmação; não restaura tabelas nem desfaz avisos já enviados. O seek temporal segue tempo de publicação, que pode diferir do tempo de negócio. Prevê ainda uma transição de entrega depois do pedido. Para comparar uma correção de cálculo, escolhe um destino isolado e controla explicitamente efeitos externos. Uma chamada de API para notificar um cliente não fica idempotente apenas porque o processamento interno usa exactly-once. Constrói um manifesto de recuperação: motivo, versão do código, intervalo pretendido, dados disponíveis, chave de deduplicação, saída de comparação e critério de promoção. Na aplicação, separa falhas transitórias de dados inválidos. Uma mensagem inválida não se torna válida por repetir a mesma transformação indefinidamente. Implementa encaminhamento de falhas no pipeline quando necessário; um dead-letter topic na subscrição não cobre automaticamente erros depois da confirmação na primeira etapa. A recuperação deve permitir corrigir e voltar a processar os registos selecionados, mantendo a ligação ao erro inicial. No comité de mudança, apresenta diferenças por chave e período, contagens de rejeições e efeitos bloqueados durante o ensaio. Um pedido de launch aceite é apenas o início desta evidência.

6. Laboratório: um contrato pequeno e observável

O ficheiro content/labs/pde-event-replay/run.py implementa um modelo original em Python, sem rede nem dependências cloud. Executa python3 content/labs/pde-event-replay/run.py < content/labs/pde-event-replay/case.json a partir da raiz do projeto. O input tem windowSeconds, allowedLatenessSeconds e uma sequência events. Cada operação é um evento com id, eventTime e value inteiros nos campos numéricos, ou um avanço explícito de watermark. Os tempos são segundos não negativos numa escala fictícia. As janelas são semiabertas: com largura 60, o instante 60 pertence à janela [60,120). O programa rejeita campos inesperados, IDs inválidos, watermark regressivo e números fora dos limites. A política local é deliberadamente explícita. Para um ID novo, watermark >= fim da janela + tolerância resulta em too-late. Antes dessa fronteira, o evento é aceite; se watermark >= fim, recebe accepted-late. Para IDs já aceites, conteúdo igual resulta em duplicate e conteúdo diferente em conflict, mesmo após o fecho. O primeiro conteúdo aceite fica conservado. O modelo guarda os IDs durante esta execução limitada a 10 000 operações e não os persiste entre processos. Não implementa triggers, panes, estado distribuído, checkpoints nem os detalhes de um runner Beam. Serve para discutir um contrato; não simula nem valida Dataflow.

7. Prever o resultado e procurar contraexemplos

Antes de executar o exemplo, segue os nove passos manualmente. A(10,7) é aceite e a sua repetição é duplicate. O watermark avança para 70. B(20,5) entra como accepted-late na primeira janela, ainda dentro da tolerância 30. A(10,9) é conflict e não substitui 7. O watermark passa a 90. C(30,4) é too-late na fronteira inclusiva. D(60,3) pertence à segunda janela e é aceite. A última repetição de A continua a ser duplicate porque a memória de IDs permanece. Espera soma 12 e contagem 2 na primeira janela, soma 3 e contagem 1 na segunda. O fecho por política só se aplica à primeira. Agora altera uma condição de cada vez. Usa watermark 89 antes de C: o evento passa a accepted-late. Usa tolerância zero: no watermark 60 já não se aceita um ID novo da primeira janela. Troca a ordem dos conteúdos conflitantes de A: o valor conservado muda, pois a regra é primeiro aceite. Repete todos os eventos num processo novo: a memória anterior desapareceu. Explica por escrito por que razão estes resultados não demonstram exatamente uma escrita num destino externo. Por fim, apresenta a evidência que faltaria para autorizar um replay real: persistência, controlo de concorrência, validade da origem e comportamento do destino perante repetições.

8. Entregar à operação com critérios verificáveis

O runbook deve ligar cada sintoma a uma decisão. Backlog a crescer pede investigação de capacidade e dependências; conflitos pedem análise de identidade e conteúdo; aumento de too-late pede confronto entre disponibilidade da origem e política temporal. Define responsáveis, evidência mínima e limites de intervenção. Para uma mudança planeada, descreve o efeito de parar a leitura e terminar trabalho pendente. Em Dataflow, drain pode emitir resultados de janelas ainda incompletas face ao período de negócio. Cancel pode deixar dados em trânsito por reconciliar. Regista a fronteira entre execuções e a forma como o consumidor trata resultados emitidos durante a paragem. Na reunião de passagem a RUN, pede ao colega que explique como detetaria um relatório incompleto apesar de o job terminar com sucesso. Dá-lhe um caso de publicação repetida, outro de conteúdo conflitante e outro de atraso além da tolerância. A resposta deve identificar a evidência disponível, a decisão permitida e a validação posterior. Resume o percurso num registo de aceitação: versão executada, dados abrangidos, registos aceites e excluídos, diferenças, efeitos externos e aprovação. A lição central é que a garantia técnica precisa de âmbito. Identidade, tempo e destino devem ter contratos compatíveis para que um resultado de processamento possa apoiar uma decisão de negócio.

"""Original bounded teaching model. No Beam runner, network or durable state."""
import json
import re
import sys


def integer(value, low, high):
    if type(value) is not int or not low <= value <= high:
        raise ValueError('integer outside the declared contract')
    return value


def exact(value, keys):
    if type(value) is not dict or set(value) != set(keys):
        raise ValueError('unexpected or missing fields')


def evaluate(document):
    exact(document, ['windowSeconds', 'allowedLatenessSeconds', 'events'])
    width = integer(document['windowSeconds'], 1, 86400)
    grace = integer(document['allowedLatenessSeconds'], 0, 86400)
    events = document['events']
    if type(events) is not list or len(events) > 10000:
        raise ValueError('events must be a list of at most 10000 operations')
    # Validate the whole input before producing a result.
    last_watermark = 0
    for event in events:
        if type(event) is dict and set(event) == {'watermark'}:
            last_watermark = integer(event['watermark'], last_watermark, 10**9)
        else:
            exact(event, ['id', 'eventTime', 'value'])
            if type(event['id']) is not str or not re.fullmatch(r'[A-Za-z0-9_-]{1,64}', event['id']):
                raise ValueError('invalid logical event id')
            integer(event['eventTime'], 0, 10**9)
            integer(event['value'], -10**6, 10**6)
    watermark = 0
    accepted = {}
    decisions = []
    for index, event in enumerate(events):
        if 'watermark' in event:
            watermark = event['watermark']
            decisions.append({'index': index, 'outcome': 'watermark', 'value': watermark})
            continue
        identifier = event['id']
        payload = (event['eventTime'], event['value'])
        end = (event['eventTime'] // width + 1) * width
        if identifier in accepted:
            outcome = 'duplicate' if accepted[identifier] == payload else 'conflict'
        elif watermark >= end + grace:
            outcome = 'too-late'
        else:
            accepted[identifier] = payload
            outcome = 'accepted-late' if watermark >= end else 'accepted'
        decisions.append({'index': index, 'id': identifier, 'outcome': outcome})
    windows = {}
    for event_time, value in accepted.values():
        start = event_time // width * width
        row = windows.setdefault(start, {'start': start, 'end': start + width, 'count': 0, 'sum': 0})
        row['count'] += 1
        row['sum'] += value
    for row in windows.values():
        row['closedByLocalPolicy'] = watermark >= row['end'] + grace
    return {'accepted': [{'id': key, 'eventTime': value[0], 'value': value[1]} for key, value in sorted(accepted.items())],
            'windows': [windows[key] for key in sorted(windows)], 'decisions': decisions,
            'watermark': watermark, 'stateLifetime': 'this invocation only',
            'sourceCompletenessProven': False, 'cloudSemanticsValidated': False,
            'productionReplayApproved': False}


def unique_object(pairs):
    result = {}
    for key, value in pairs:
        if key in result:
            raise ValueError('duplicate JSON field')
        result[key] = value
    return result


if __name__ == '__main__':
    try:
        raw = sys.stdin.read(2_000_001)
        if len(raw) > 2_000_000:
            raise ValueError('input exceeds local limit')
        result = evaluate(json.loads(raw, object_pairs_hook=unique_object))
        print(json.dumps(result, sort_keys=True))
    except (ValueError, TypeError, RecursionError) as error:
        print(json.dumps({'error': str(error)}), file=sys.stderr)
        sys.exit(2)
NA PRÁTICA

A(10,7) e B(20,5) somam 12; A(10,9) é conflito e C após a fronteira é too-late.

Armadilhas comuns

Confundir mensagem com evento; substituir conflitos; inferir completude do watermark; repetir avisos durante replay.

Tópicos relacionados: Reconciliação de migrações · Modelação de dados · Observabilidade e SLO

Leva esta ideia contigo

Especifica identidade, política temporal e efeitos no destino; reconcilia o resultado antes de o aceitar.

Criar conta

Referência: Professional Data Engineer standard exam guide · Current linked standard guide (document title v4.2); edition date unconfirmed (2026-09-30 inspection)

Google Cloud é uma marca comercial de Google LLC. A dr.pt é uma plataforma de preparação independente e não está afiliada, associada, patrocinada, autorizada nem aprovada por Google. Os conteúdos e as perguntas são originais, não são perguntas oficiais de exame, e concluir os nossos testes não atribui nem garante qualquer certificação. Os nomes são usados apenas para identificar o tema. Todas as outras marcas pertencem aos respetivos titulares.