Definir o que tem de chegar e quando
Uma migração precisa de uma definição verificável de conteúdo, não apenas de uma data no calendário. Especifica tabelas ou ficheiros, filtros, histórico, alterações posteriores e regras de transformação. Num projeto fictício de reporting de fundos, o responsável de dados pode exigir todas as operações aceites até às 22:00, incluindo o batch que começou antes dessa hora. Essa frase obriga a identificar produtores e a demonstrar o corte. Dimensiona também a transferência: 10 TB decimais a 500 Mbit/s efetivos exigem pelo menos 160.000 segundos, cerca de 44,4 horas, antes de validação e margem. Se entram 25 MB/s de alterações e só aplicas 20 MB/s, o backlog aumenta 5 MB/s. O ensaio deve medir extração, transporte e aplicação; aumentar apenas a ligação pode deixar o verdadeiro gargalo intacto. Documenta unidades e pressupostos para que o steering committee possa distinguir uma previsão medida de uma esperança de conclusão.
Comparar estados com uma fronteira demonstrável
A cópia inicial representa uma base para a migração; as alterações seguintes têm de ser capturadas pelo mecanismo escolhido. Uma consulta feita às 10:00 na origem e outra às 10:05 no destino podem observar estados legitimamente diferentes. Isso não autoriza classificar todas as divergências como lag, nem como corrupção. Conserva os resultados e estabelece um corte comparável através das capacidades dos sistemas e de um procedimento ensaiado. No caso ativo-passivo desta aula, o destino fica sem escritas de negócio até controlar produtores, terminar trabalho aceite, aplicar alterações e reconciliar. Um ficheiro CFT em processamento faz parte dessa fronteira, mesmo que a interface web já esteja parada. Depois de aceitar escritas no destino, a origem antiga deixa de ser um fallback automaticamente equivalente. Regista como serão preservadas essas novas operações antes de prometer retorno. Um nome como cut-9 num relatório é apenas um rótulo; precisa de evidência que o ligue ao estado observado.
Preservar o significado durante a transformação
Escreve um contrato de comparação antes de normalizar valores. No exercício, os IDs são texto: 001 e 1 são diferentes. Uma referência ausente e uma referência recebida vazia também são diferentes; não as convertas ambas para texto vazio para fazer o teste passar. Os timestamps representam instantes, por isso 11:00+01:00 e 10:00Z são equivalentes depois de interpretar o offset. Uma hora sem fuso é rejeitada, porque não existe informação suficiente para escolher o instante. Os montantes desta simulação usam exatamente duas casas decimais, uma hipótese local que não se aplica a todas as moedas. Na arquitetura BigQuery, escolhe um tipo decimal com precisão e escala adequadas e testa os limites. SAFE_CAST pode transformar uma entrada inválida em NULL; regista rejeições separadamente em vez de confiar numa soma que só observa valores válidos. Cada regra de limpeza deve ter uma justificação de negócio e uma versão: remover espaços ou ignorar maiúsculas pode alterar identificadores que parecem apenas texto.
Combinar controlos que observam falhas diferentes
Um checksum de transporte verifica integridade dos bytes; não confirma que a transformação posterior conserva o significado. Contagem, soma, unicidade e comparação por identidade respondem a perguntas diferentes. Se A passa de 100 para 110 e B de 200 para 190, a contagem e a soma mantêm-se. Se um ID desaparece e outro surge com o mesmo valor, até os valores agregados por moeda podem coincidir. Começa por verificar âmbito e chaves, compara valores relevantes e usa agregados como evidência complementar. No Data Validation Tool, escolhe uma chave que identifique uma linha do resultado; trade_id pode precisar de leg_id. Inspeciona as normalizações das consultas geradas. Um --dry-run apresenta SQL e não demonstra igualdade dos dados. Para um conjunto grande, planeia cobertura por partições ou intervalos, conservando quais foram efetivamente verificados. Uma amostra sem divergências não prova que todo o conjunto esteja correto.
Exercício: encontrar diferenças que o total esconde
Antes de executar o código, prevê três resultados: trocar a ordem das linhas, substituir 001 por 1 e redistribuir um euro entre dois IDs mantendo a soma. Executa python3 content/labs/pca-data-reconciliation/run.py. Os controlos observados devem mostrar equivalência no primeiro caso e divergências nos outros dois. A função canonical valida o pequeno esquema, converte montantes decimais para cêntimos inteiros e normaliza instantes para UTC com precisão até microssegundos. A função index rejeita IDs repetidos; não elimina silenciosamente linhas. compare identifica IDs em falta, inesperados e campos alterados, e conserva totais separados por moeda. Modifica uma referência de None para texto vazio e verifica o campo reportado. O programa só compara listas sintéticas em memória; não executa DVT, não captura snapshots nem contacta Google Cloud. Duas listas vazias com o mesmo rótulo passam, mas isso não prova extração completa. No trabalho real, guarda também a evidência do âmbito e do corte, com acesso adequado aos resultados.
Desenhar o significado de um resultado em streaming
Uma atualização mais recente na origem pode chegar antes de uma antiga a um consumidor próprio de ficheiros Datastream. A ordem de chegada não deve substituir a metadata de sequência aplicável à origem. Não atribuas ao teu consumidor o merge que um destino gerido pode implementar. Se depois agregas eventos, decide qual timestamp representa o negócio: um pagamento às 20:59 que chega às 21:02 pertence à janela de ocorrência, quando esse é o contrato. O watermark ajuda a decidir quando emitir, mas ainda podem chegar dados atrasados. Define o que pode ser corrigido e durante quanto tempo. Com panes acumulados, emitir 100 e depois 130 para a mesma janela não significa produzir 230. O destino precisa de reconhecer revisões e evitar que uma revisão antiga substitua a nova. Estes detalhes afetam o dashboard que a equipa de operações usa para decidir; um resultado rápido sem significado de completude pode induzir uma decisão errada.
Tornar consultas e estimativas compatíveis com o uso
Escolhe organização dos dados a partir das consultas reais. Se o reporting filtra data de negócio e mesa, ensaia particionamento pela data e clustering por mesa dentro das partições. Mede o efeito com dados representativos; o nome da funcionalidade não garante uma poupança. Numa tabela nativa não clustered, SELECT * LIMIT 10 não substitui o filtro da partição nem a seleção de colunas. Confirma que o predicado permite pruning e observa bytes previstos e reais. Também não interpretes zero bytes num dry run de uma tabela externa como promessa de custo zero: pode ser apenas o limite inferior que o estimador consegue fornecer. Se declaras primary keys em BigQuery, mantém os dados coerentes com elas, porque não são impostas automaticamente e podem influenciar otimizações. A arquitetura deve incluir quem valida essas propriedades quando chegam novos dados, não apenas quem criou a tabela durante o projeto.
Fechar a decisão com critérios e responsáveis
Prepara uma tabela de aceitação com requisito, observação, resultado, responsável e ação pendente. No caso dos três trades, escreve explicitamente que o total coincide mas T2 diverge, T3 falta e T9 é inesperado. A decisão segue a regra aprovada de equivalência por ID, portanto a aceitação fica bloqueada. Uma exceção exige avaliar impacto e autoridade de aprovação; não deve nascer de alterar o verificador para obter verde. Para a passagem a RUN, inclui como repetir a reconciliação, onde observar rejeições, quando um relatório é provisório e como escalar uma diferença. Se o batch diário já cumpre a necessidade medida, usa-o como referência e pede um benefício concreto antes de acrescentar streaming. Resume o princípio da aula: demonstrar transporte, estado, significado e utilização requer evidências diferentes. O arquiteto e o gestor de projeto devem ligar essas evidências à decisão, preservando as limitações que ainda não foram resolvidas.
"""Original local reconciliation exercise. Synthetic data; no cloud or DVT calls.
Contract: unique string IDs; exact two-decimal amounts; aware instants;
case-sensitive references including NULL versus empty; same declared checkpoint.
The checkpoint is an assertion by the caller, not proof of a consistent snapshot.
"""
import json
import re
from datetime import datetime, timezone
from decimal import Decimal
FIELDS = {'id', 'currency', 'amount', 'occurred_at', 'reference'}
def canonical(row):
if not isinstance(row, dict) or set(row) != FIELDS:
raise ValueError('exact exercise schema required')
if not isinstance(row['id'], str) or not row['id']:
raise ValueError('nonempty string ID required')
if not isinstance(row['currency'], str) or not re.fullmatch(r'[A-Z]{3}', row['currency']):
raise ValueError('three uppercase letters required; not an ISO currency registry check')
amount = row['amount']
if not isinstance(amount, str) or not re.fullmatch(r'[+-]?\d{1,12}(?:\.\d{1,2})?', amount, flags=re.ASCII):
raise ValueError('bounded exact decimal string required')
cents = int(Decimal(amount) * 100)
instant = row['occurred_at']
if not isinstance(instant, str) or not re.fullmatch(r'\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}(?:\.\d{1,6})?(?:Z|[+-]\d{2}:\d{2})', instant, flags=re.ASCII):
raise ValueError('timestamp string required')
try:
parsed = datetime.fromisoformat(instant)
except ValueError as exc:
raise ValueError('invalid ISO timestamp') from exc
if parsed.tzinfo is None or parsed.utcoffset() is None:
raise ValueError('explicit UTC offset required')
ref = row['reference']
if ref is not None and not isinstance(ref, str):
raise ValueError('reference must be a string or None')
return {'id': row['id'], 'currency': row['currency'], 'cents': cents,
'instant': parsed.astimezone(timezone.utc).isoformat(timespec='microseconds'),
'reference': ref}
def index(rows):
result = {}
for row in rows:
item = canonical(row)
if item['id'] in result:
raise ValueError('duplicate ID: ' + item['id'])
result[item['id']] = item
return result
def totals(items):
result = {}
for row in items.values():
bucket = result.setdefault(row['currency'], {'count': 0, 'cents': 0})
bucket['count'] += 1
bucket['cents'] += row['cents']
return result
def compare(source, target, source_checkpoint, target_checkpoint):
if not isinstance(source_checkpoint, str) or not source_checkpoint or source_checkpoint != target_checkpoint:
raise ValueError('same nonempty declared checkpoint required')
left, right = index(source), index(target)
missing = sorted(left.keys() - right.keys())
unexpected = sorted(right.keys() - left.keys())
changed = {key: sorted(k for k in left[key] if left[key][k] != right[key][k])
for key in sorted(left.keys() & right.keys()) if left[key] != right[key]}
return {'same': not (missing or unexpected or changed), 'missing': missing,
'unexpected': unexpected, 'changed': changed,
'sourceTotals': totals(left), 'targetTotals': totals(right)}
def sample(key='001', amount='12.30', **changes):
return dict(id=key, currency='EUR', amount=amount,
occurred_at='2026-10-05T10:00:00+00:00', reference=None) | changes
def run():
checks = []
def check(name, condition):
if not condition:
raise AssertionError(name)
checks.append(name)
def reject(name, fn):
try:
fn()
except ValueError:
checks.append(name)
else:
raise AssertionError(name)
def same(a, b): return compare(a, b, 'cut-9', 'cut-9')
a = sample()
check('decimal amount preserved', canonical(a)['cents'] == 1230)
check('negative exact amount', canonical(sample(amount='-0.01'))['cents'] == -1)
check('equivalent decimal spelling', canonical(sample(amount='12.3')) == canonical(a))
check('same instant different offset', canonical(sample(occurred_at='2026-10-05T11:00:00+01:00')) == canonical(a))
check('UTC Z accepted', canonical(sample(occurred_at='2026-10-05T10:00:00Z')) == canonical(a))
check('identity retained as text', canonical(a)['id'] == '001')
check('leading zeros remain significant', same([a], [sample(key='1')])['missing'] == ['001'])
check('NULL differs from empty', same([a], [sample(reference='')])['changed'] == {'001': ['reference']})
check('reference case significant', not same([sample(reference='ab')], [sample(reference='AB')])['same'])
check('reference whitespace significant', not same([sample(reference='ab')], [sample(reference='ab ')])['same'])
check('same rows reordered', same([a, sample('002')], [sample('002'), a])['same'])
check('missing row identified', same([a, sample('002')], [a])['missing'] == ['002'])
check('unexpected row identified', same([a], [a, sample('003')])['unexpected'] == ['003'])
left = [sample('001', '10.00'), sample('002', '20.00')]
right = [sample('001', '11.00'), sample('002', '19.00')]
diff = same(left, right)
check('equal count and sum hide changes', diff['sourceTotals'] == diff['targetTotals'] and not diff['same'])
check('both cancelling changes identified', diff['changed'] == {'001': ['cents'], '002': ['cents']})
replacement = same([sample('001', '10.00')], [sample('009', '10.00')])
check('replacement identity not equivalent', replacement['missing'] == ['001'] and replacement['unexpected'] == ['009'])
mixed = index([sample('001', '10.00'), sample('002', '10.00', currency='USD')])
check('currencies never added together', totals(mixed) == {'EUR': {'count': 1, 'cents': 1000}, 'USD': {'count': 1, 'cents': 1000}})
check('empty pair equivalent for supplied scope', same([], [])['same'])
check('microseconds retained', not same([a], [sample(occurred_at='2026-10-05T10:00:00.000001Z')])['same'])
check('Unicode preserved', same([sample(reference='ação')], [sample(reference='ação')])['same'])
reject('duplicate source rejected', lambda: same([a, a], [a]))
reject('duplicate target rejected', lambda: same([a], [a, a]))
reject('different checkpoints rejected', lambda: compare([a], [a], 'cut-9', 'cut-10'))
reject('empty checkpoint rejected', lambda: compare([], [], '', ''))
reject('float amount rejected', lambda: canonical(sample(amount=12.3)))
reject('extra precision rejected', lambda: canonical(sample(amount='12.301')))
reject('NaN rejected', lambda: canonical(sample(amount='NaN')))
reject('missing offset rejected', lambda: canonical(sample(occurred_at='2026-10-05T10:00:00')))
reject('invalid timestamp rejected', lambda: canonical(sample(occurred_at='2026-99-05T10:00:00Z')))
reject('missing field rejected', lambda: canonical({k: v for k, v in a.items() if k != 'reference'}))
reject('numeric ID rejected', lambda: canonical(sample(key=1)))
reject('lowercase currency rejected', lambda: canonical(sample(currency='eur')))
reject('oversized decimal rejected', lambda: canonical(sample(amount='1000000000000.00')))
reject('nontext reference rejected', lambda: canonical(sample(reference=0)))
reject('extra field rejected', lambda: canonical(a | {'other': 1}))
reject('empty ID rejected', lambda: canonical(sample(key='')))
check('JSON retains null versus string', json.loads(json.dumps([None, ''])) == [None, ''])
reject('submicrosecond input rejected', lambda: canonical(sample(occurred_at='2026-10-05T10:00:00.0000001Z')))
check('input records unchanged', a == sample())
return {'passed': len(checks), 'checks': checks, 'example': diff,
'network': False, 'vendorExecution': False, 'persistentWrites': False}
if __name__ == '__main__':
print(json.dumps(run(), ensure_ascii=False, indent=2))
Um lote com três trades conserva 600 EUR, mas muda um valor e substitui um ID. Contagem e soma passam enquanto a aceitação por registo falha.
Armadilhas comuns
Confundir checksum com equivalência semântica, totais iguais com linhas iguais, cópia inicial com sincronização final e rótulos de corte com snapshots demonstrados.
Tópicos relacionados: Dados, consistência e publicação de eventos · Migração, custos e aceitação · Observabilidade e evidência de conclusão
Define o âmbito e a fronteira; preserva identidade e significado; combina controlos independentes e liga a decisão de aceitação à evidência produzida.
Referência: Transfer your large datasets · Current linked standard guide; edition date unconfirmed (2026-09-30 inspection)