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

Recuperar a ingestão e reconciliar alterações

Retomar CDC e snapshots com controlo de versões, eliminações, intervalos perdidos e evidência de reconciliação.

1. Definir o que tem de ser recuperado

Uma falha de ingestão deixa pelo menos três perguntas em aberto: que alterações ainda existem na origem, que estado o pipeline guardou e que efeitos já chegaram ao destino. Um processo que volta a executar não responde automaticamente às três. Antes de escolher replay ou backfill, regista o último ponto confirmado em cada fronteira e a evidência que o sustenta. Distingue um timestamp de observação de uma posição de replicação: o primeiro diz quando olhaste; a segunda participa no contrato de continuidade da origem. Num exemplo fictício de posições de fundos, a equipa necessita dos saldos atuais para uma consulta operacional e de todas as alterações intradiárias para um controlo separado. Um backfill pode ajudar a recompor os saldos sem recuperar cada operação intermédia. Cria dois critérios de aceitação e dois responsáveis. Se a história não puder ser reconstruída, mantém uma exceção explícita para esse requisito. Igualdade dos saldos não autoriza declarar que nenhuma operação se perdeu. Define também quem pode suspender a passagem e onde ficam os resultados da reconciliação, evitando que a pressão do fecho transforme uma hipótese em evidência.

2. Retomar a origem sem esconder um intervalo perdido

No Datastream, recuperar um stream permanentemente falhado exige escolher uma posição adequada ao tipo de origem. Se um binlog MySQL necessário pode ser recuperado, avalia repor esse ficheiro e tentar a posição atual antes de saltar para a mais recente. Num slot PostgreSQL novo, as alterações entre a posição perdida e o primeiro LSN disponível não reaparecem por o stream voltar a funcionar. Regista esse intervalo como uma obrigação de recuperação ou reconciliação, com as limitações do histórico disponíveis. Não compares posições de duas origens apenas porque os nomes de base de dados coincidem. Depois de um failover, verifica a correspondência da sequência no servidor que passou a produzir alterações. Um teste de ligação confirma conectividade, mas não demonstra continuidade do log. Para planear a intervenção, anota a retenção disponível, os objetos afetados, as escritas que continuam durante a mudança e o destino onde ensaiar a sobreposição. A decisão sobre a posição inicial deve ser explicável a quem validar os dados depois. Uma recuperação rápida que salta um intervalo precisa de um plano explícito para esse intervalo.

3. Controlar a leitura histórica e a aplicação

O estado Completed de um backfill confirma que o objeto foi lido; o carregamento no destino pode continuar. Por isso, o responsável pela passagem não deve aceitar apenas uma captura de ecrã do estado. Pede evidência do destino num corte comparável, incluindo chaves, valores e eliminações relevantes. Quando houver escrita concorrente, documenta como os eventos de CDC e as imagens históricas são ordenados. O último ficheiro recebido pode conter uma imagem mais antiga do que o estado já aplicado. Parar um backfill pode ser uma medida válida para proteger a origem de carga excessiva. Contudo, ao reiniciar um objeto parado, o plano deve prever nova transmissão dos seus dados existentes. O orçamento de tempo e carga não deve assumir continuação exata na linha onde a leitura parou. Ensaiar um objeto de dimensão representativa ajuda a escolher concorrência e janela. Não extrapoles para todas as tabelas a partir de uma tabela pequena e sem alterações. Regista throughput, impacto na origem, atraso de aplicação e critérios para reduzir a concorrência ou interromper a operação.

4. Separar a ordem do produto da política do exercício

Datastream fornece metadata de eventos, incluindo identidade e informação de ordenação dependente da origem. Um consumidor próprio não deve deduplicar todas as alterações de uma chave como se fossem uma única entrega. Uma atualização e uma eliminação da mesma linha são eventos legítimos diferentes. No BigQuery CDC, a sequência personalizada permite ordenar alterações da mesma chave; as secções hexadecimais são comparadas numericamente e por ordem. Se as sequências forem iguais, o tempo de ingestão desempata. Misturar escritas com e sem a sequência personalizada não constitui uma política previsível. O laboratório abaixo usa uma regra diferente e deliberadamente mais conservadora. Cada chave tem uma versão inteira inventada para o exercício, sem pretender representar LSN, SCN ou o formato BigQuery. Conteúdos diferentes para a mesma chave e versão bloqueiam a chave, mesmo quando o conflito é anterior à maior versão. Esta escolha torna a divergência visível para investigação. Não a apresentes como uma garantia do fornecedor. Ao adaptar o raciocínio a um pipeline, define quem produz a versão, como se compara e o que acontece quando duas fontes discordam.

5. Executar a reconciliação local

O exercício recebe baseline, events e expected. Cada linha da baseline e da referência tem key, version, deleted e value. Uma eliminação conserva a versão e usa value=null; uma linha ativa usa um inteiro. Os eventos acrescentam id e op, com upsert ou delete. Há um máximo de mil elementos por lista. O programa valida a entrada completa, rejeita chaves repetidas nas listas de estados e rejeita reutilização do mesmo id de evento com conteúdo diferente. Duas entregas idênticas contam como uma entrega duplicada, sem criar uma nova alteração. No ficheiro case.json, A começa na versão 4 com valor 100; B está eliminado na versão 7. Chegam versões 6 e 5 de A, uma versão 6 de B, uma nova chave C e uma repetição de e1. Executa python3 content/labs/pde-cdc-reconcile/run.py < content/labs/pde-cdc-reconcile/case.json. A termina em 120 na versão 6, B continua eliminado e C vale 30. O resultado identifica uma entrega duplicada e coincide com a referência declarada. Essa referência é fornecida pelo utilizador do exercício; o programa não confirma que corresponde à origem real.

6. Tentar refutar o resultado

Muda o valor de A na referência para 80 e o de C para 70. A soma continua em 150, mas a comparação por chave deve produzir duas divergências. Depois, acrescenta outro evento para A, versão 5, com valor 111: a chave fica bloqueada por conflito, apesar de existir uma versão 6. Finalmente, acrescenta uma versão 8 de B com valor 55. Neste contrato, uma alteração posterior pode voltar a ativar a linha; conservar um tombstone não significa impedir para sempre qualquer nova versão. Troca a ordem dos eventos e confirma que os estados e conflitos permanecem iguais. Essa propriedade resulta da seleção por versão e da recolha de conflitos, não da ordem de chegada. Retira uma chave da referência para observar unexpected-state; acrescenta uma chave inexistente para obter missing-state. O resultado também compara a versão, mesmo quando o valor coincide. Nenhum destes testes prova atomicidade de uma transação entre várias chaves, durabilidade após reinício ou completude de todo o histórico. O estado existe apenas durante esta invocação. Esses limites devem acompanhar qualquer conclusão apresentada no ensaio.

7. Coordenar snapshot, replay e destino

Um snapshot Dataflow pode guardar estado do job e, quando configurado, snapshots das fontes Pub/Sub. Não inclui um snapshot automático do sink. Se a tabela recebeu resultados depois do ponto guardado, restaurar o pipeline não desfaz esses resultados. O plano tem de controlar a sobreposição: por exemplo, através de uma escrita idempotente comprovada para o contrato ou de um destino isolado com reconciliação antes de promoção. Não declares esta garantia apenas porque a entrega de mensagens tem uma opção chamada exactly-once. Verifica região, compatibilidade da nova versão e validade de cada snapshot associado. A validade da fonte pode terminar antes da validade do estado Dataflow. Um seek também tem consequências no controlo de entrega: o filtro da subscrição que o executa limita a redelivery e a contagem de tentativas de dead-letter é reiniciada. Mantém evidência histórica fora desses contadores. No ensaio, explica quais os consumidores que continuam ativos, como se distingue trabalho anterior de trabalho repetido e como se reconhece uma confirmação antiga que já não é válida depois do seek.

8. Entregar uma decisão verificável ao RUN

A passagem deve ligar o incidente a um conjunto pequeno de evidências verificáveis: intervalo afetado, posição de reinício, versão do pipeline, recursos usados no restauro e diferenças encontradas no destino. Num contexto internacional, prepara uma nota curta em inglês com o estado confirmado, a incerteza restante e o próximo responsável. Evita escrever recovered sem dizer se isso significa processo ativo, estado atual reconciliado ou histórico completo. Cada significado exige prova diferente e pode ter uma hora de conclusão distinta. No BigQuery CDC, a tolerância max_staleness pode permitir uma leitura de baseline anterior. Escolhe um corte comparável para a reconciliação e verifica a aplicação das alterações antes de interpretar diferenças como perda. A revisão final deve incluir eliminações, conflitos, chaves inesperadas e valores que se compensam em totais globais. Um relatório verde do laboratório só significa correspondência com a referência fornecida e ausência dos conflitos modelados. Resume assim a aprendizagem: recuperar a leitura, controlar a ordem, preservar a evidência de eliminações, reconciliar o destino e separar estado atual de história. A autorização de produção continua a pertencer ao processo real da organização. Uma nota de passagem pode indicar: estado atual reconciliado no corte acordado; histórico intradiário ainda sob investigação; replay limitado ao destino de ensaio; próximo responsável identificado. Acrescenta links para as consultas e os resultados, a versão dos scripts e a hora de recolha. Outra pessoa deve conseguir repetir a comparação sem depender da memória de quem tratou o incidente. Se uma premissa mudar, como a retenção disponível ou a continuidade da origem, revê a decisão e regista a nova evidência. Este encadeamento ajuda a separar uma mitigação temporária de uma recuperação aceite pelo negócio.

"""Original bounded versioned-row teaching model; no cloud or durable writes."""
import json
import re
import sys


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


def identifier(value):
    if type(value) is not str or not re.fullmatch(r'[A-Za-z0-9_-]{1,48}', value):
        raise ValueError('invalid identifier')
    return value


def integer(value, lo, hi):
    if type(value) is not int or not lo <= value <= hi:
        raise ValueError('integer outside contract')
    return value


def row(value):
    exact(value, ['key', 'version', 'deleted', 'value'])
    identifier(value['key'])
    integer(value['version'], 0, 10**9)
    if type(value['deleted']) is not bool:
        raise ValueError('deleted must be boolean')
    if value['deleted']:
        if value['value'] is not None:
            raise ValueError('tombstone value must be null')
    else:
        integer(value['value'], -10**9, 10**9)
    return dict(value)


def row_list(values):
    if type(values) is not list or len(values) > 1000:
        raise ValueError('at most 1000 rows')
    result = {}
    for value in values:
        r = row(value)
        if r['key'] in result:
            raise ValueError('duplicate row key')
        result[r['key']] = r
    return result


def evaluate(document):
    exact(document, ['baseline', 'events', 'expected'])
    baseline, expected = row_list(document['baseline']), row_list(document['expected'])
    events = document['events']
    if type(events) is not list or len(events) > 1000:
        raise ValueError('at most 1000 events')
    unique, duplicate_deliveries = {}, 0
    for event in events:
        exact(event, ['id', 'key', 'version', 'op', 'value'])
        identifier(event['id'])
        if event['op'] not in ('upsert', 'delete'):
            raise ValueError('unknown operation')
        r = row({'key': event['key'], 'version': event['version'],
                 'deleted': event['op'] == 'delete', 'value': event['value']})
        if event['id'] in unique:
            if unique[event['id']] != r:
                raise ValueError('event identity reused with different content')
            duplicate_deliveries += 1
        unique[event['id']] = r
    versions = {}
    for r in list(baseline.values()) + list(unique.values()):
        versions.setdefault(r['key'], {}).setdefault(r['version'], set()).add((r['deleted'], r['value']))
    states, conflicts = {}, []
    for key, history in sorted(versions.items()):
        bad = sorted(v for v, values in history.items() if len(values) > 1)
        if bad:
            conflicts.append({'key': key, 'versions': bad})
            continue
        version = max(history)
        deleted, value = next(iter(history[version]))
        states[key] = {'key': key, 'version': version, 'deleted': deleted, 'value': value}
    blocked = {r['key'] for r in conflicts}
    differences = []
    for key in sorted(set(states) | set(expected) | blocked):
        if key in blocked:
            reason = 'conflicting-version'
        elif key not in states:
            reason = 'missing-state'
        elif key not in expected:
            reason = 'unexpected-state'
        elif states[key] != expected[key]:
            reason = 'state-mismatch'
        else:
            continue
        differences.append({'key': key, 'reason': reason})
    return {'states': [states[k] for k in sorted(states)], 'conflicts': conflicts,
            'differences': differences, 'duplicateDeliveries': duplicate_deliveries,
            'matchesDeclaredReference': not differences,
            'referenceIndependentlyVerified': False, 'historyCompletenessProven': False,
            'cloudSemanticsValidated': False, 'productionReplayApproved': False,
            'stateLifetime': 'this invocation only'}


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')
        print(json.dumps(evaluate(json.loads(raw, object_pairs_hook=unique_object)), sort_keys=True))
    except (ValueError, TypeError, RecursionError) as error:
        print(json.dumps({'error': str(error)}), file=sys.stderr)
        sys.exit(2)
NA PRÁTICA

Caso fictício: o backfill repõe saldos, mas a equipa de controlo ainda exige operações intradiárias. A decisão distingue estado atual, histórico e passagem ao RUN.

Armadilhas comuns

Confundir Completed com dados aplicados; usar chegada como versão; eliminar tombstones demasiado cedo; aceitar apenas totais; presumir que um snapshot Dataflow reverte o sink.

Tópicos relacionados: Continuidade de CDC · Snapshots e retenção · Idempotência e reconciliação

Leva esta ideia contigo

O replay tem de respeitar versões e efeitos anteriores. Estado reconciliado não demonstra histórico completo nem autoriza uma passagem de produção.

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.