Definir a regra e a fronteira transacional
Começa pela regra de negócio que não pode ser violada. Uma reserva exige quantidade disponível e tem de atualizar o registo de reserva e a disponibilidade em conjunto. Na mesma base Spanner, uma transação read-write serializable pode incluir a leitura crítica e ambas as escritas. Uma leitura strong isolada não transforma pedidos seguintes numa única decisão atómica. Para uma reconciliação só de leitura com várias consultas, usa um snapshot comum através de uma transação read-only. Duas leituras independentes podem observar instantes diferentes, mesmo sendo atuais. Escolher stale reads é aceitar um instante passado: pode servir um relatório tolerante a atraso, mas não um ecrã que tem de mostrar o commit acabado de confirmar. Regista estas escolhas no desenho, ligando cada requisito à operação concreta e ao modo de isolamento, sem atribuir todas as garantias a qualquer configuração do produto.
Distinguir ordem, retry e resultado incerto
Em transações serializable, a consistência externa preserva a precedência de commits não sobrepostos: se a revisão de uma política foi confirmada antes de começar a aprovação dependente, um snapshot com a aprovação não pode omitir essa revisão. Isto não cria uma transação com um processador externo. Um callback read-write pode executar várias vezes e uma chamada externa feita dentro dele pode produzir efeitos repetidos. Também um DEADLINE_EXCEEDED numa escrita não prova rollback nem sucesso. Antes de reenviar, procura o estado associado à identidade original e aplica o procedimento de recuperação idempotente. No exercício de liquidação, dois recibos externos exigem reconciliação mesmo que a base tenha só uma linha. Não apagues evidência nem inventes uma compensação financeira; identifica responsáveis, confirma cada resultado e conserva a relação entre intenção, tentativa e efeito observado.
Distribuir carga sem perder a identidade
Uma chave de negócio estável serve para reconhecer a mesma intenção; uma chave física também influencia onde as escritas se concentram. Não confundas os dois papéis. Em Spanner, um timestamp monotónico como primeiro componente pode enviar novas inserções para a mesma faixa. Inverter a ordenação muda o extremo, não elimina o padrão. Um prefixo distribuído ou desenho de shards deve ser avaliado juntamente com consultas temporais, índices e operação. Um índice global iniciado pelo mesmo timestamp pode reproduzir a concentração. Antes de escolher, recolhe a taxa prevista, distribuição real das chaves, consultas críticas e custo de pesquisar várias faixas. Num projeto de migração, esta análise deve entrar no ensaio de carga e nos critérios de aceitação. Mais capacidade não é prova de que o modelo de acesso está bem distribuído.
Escolher a menor sequência que o negócio exige
Em Pub/Sub, a ordem é por chave, com configuração de ordering na subscrição e publicação da mesma chave na mesma região. Não é uma ordem global entre contas diferentes. Se o negócio só exige sequência por conta, uma chave única para todo o banco cria serialização desnecessária. Usa granularidade compatível com a regra e mede o backlog. Ordenação também não elimina reentrega: numa sequência at-least-once sem dead-letter topic, uma mensagem reentregue pode trazer de novo mensagens posteriores da chave, mesmo acknowledged. O consumidor precisa de reconhecer efeitos já aplicados. Se o callback agenda trabalho assíncrono e regressa, a ordem de receção não garante ordem de conclusão desse trabalho. Desenha explicitamente a sequência dos efeitos e o momento de acknowledgment, para que uma otimização de concorrência não altere o resultado de negócio.
Preparar falhas e replay como capacidades operacionais
Um dead-letter topic não é apenas um nome na configuração. O service agent Pub/Sub precisa de publicar no destino e fazer acknowledgment na subscrição de origem. As permissões de um operador humano não substituem as dessa identidade. O número máximo de tentativas é aproximado; um requisito de exatamente cinco efeitos de negócio exige um controlo diferente. Para replay de mensagens já acknowledged, planeia retenção no tópico ou retenção dessas mensagens na subscrição antes de precisar delas. Seek altera o estado de acknowledgment, mas não desfaz débitos, emails ou outras ações já executadas. A entrega também não converge instantaneamente após seek. Prepara um procedimento que identifique o intervalo, a versão do consumidor, os critérios de conclusão e a proteção contra repetição. O ensaio deve incluir mensagens que regressam quando parte do seu trabalho já existe.
Publicar conjuntos de ficheiros por identidade de versão
Cloud Storage oferece consistência forte para leitura após escrita e listagem, mas caches públicas podem continuar a servir uma cópia antiga. Distingue a origem da camada que respondeu. Para concorrência, usa precondições: ifGenerationMatch=0 impede substituir uma versão live durante uma criação; uma eliminação deve conservar a geração aprovada para não atingir uma substituição posterior. Nas leituras por blocos, fixa a geração para evitar juntar partes de versões diferentes. Nada disto torna um batch de operações sobre vários objetos atómico. Um dataset com accounts e positions precisa de identidade de execução e gerações verificadas. Uma aplicação pode publicar um manifesto de aceitação depois de validar o conjunto; essa é uma decisão de desenho, não uma garantia automática do bucket. Confirma também quem pode alterar o manifesto e como os consumidores rejeitam conjuntos incompletos.
Separar append, finalização e visibilidade analítica
Na Storage Write API gRPC de BigQuery, o default stream tem semântica at-least-once. Um stream criado pela aplicação pode usar offsets para reconhecer appends dentro desse stream. Dois streams distintos no offset zero não representam a mesma posição, mesmo contendo o mesmo orderId. A identidade de negócio continua a precisar de tratamento. Para um lote, streams pending permitem adiar visibilidade: os workers escrevem e finalizam; o coordenador confirma o grupo quando tudo está pronto. Finalize impede novos appends, mas não publica os dados. BatchCommitWriteStreams abrange streams finalizados da mesma tabela parent. Se stream_errors não estiver vazio, nenhum stream desse grupo foi confirmado pela chamada. Conserva a resposta e o commit_time quando existe, em vez de tratar a simples receção de uma resposta como sucesso do lote inteiro.
Laboratório: comparar o estado local com o efeito externo
Executa o código abaixo com Python 3.12 ou posterior e o módulo sqlite3 disponível. Antes de executar, prevê três resultados: o que fica após o aborto de unsafe_attempt, o que acontece quando relay perde a resposta e o que muda ao criar S2 em vez de repetir S1. O laboratório usa uma base SQLite em memória e um dicionário como destinatário fictício. Demonstra atomicidade local e reconhecimento sequencial de identidade, sem executar Spanner, Pub/Sub ou pagamentos. O dicionário não é armazenamento durável nem proteção testada contra consumidores concorrentes. Num sistema real, a verificação de identidade e o efeito precisam de garantias apropriadas no destinatário. Relaciona os resultados com os dois casos finais e escreve critérios de aceitação: uma intenção recuperável, referências estáveis, inputs da mesma execução e confirmação explícita da publicação completa.
"""Original in-memory SQLite teaching lab, not a Spanner or Pub/Sub emulator.
All external calls below are Python dictionary operations. No network, credentials,
files, concurrency, vendor services or actual payments are used.
"""
import json
import platform
import sqlite3
checks = []
def check(name, condition):
if not condition:
raise AssertionError(name)
checks.append(name)
def fails(call, exception):
try:
call()
except exception:
return True
return False
db = sqlite3.connect(":memory:", isolation_level=None,
autocommit=sqlite3.LEGACY_TRANSACTION_CONTROL)
db.executescript("""
CREATE TABLE orders (id TEXT PRIMARY KEY, amount INTEGER NOT NULL);
CREATE TABLE outbox (id TEXT PRIMARY KEY, amount INTEGER NOT NULL, sent INTEGER NOT NULL);
""")
unsafe_receipts = []
def unsafe_attempt(order_id, amount, abort=False):
db.execute("BEGIN")
try:
# This fictional external list is outside SQLite's transaction boundary.
unsafe_receipts.append((order_id, amount))
db.execute("INSERT INTO orders VALUES (?, ?)", (order_id, amount))
if abort:
raise RuntimeError("injected failure before commit")
db.execute("COMMIT")
except Exception:
db.execute("ROLLBACK")
raise
def enqueue(order_id, amount, abort=False):
db.execute("BEGIN")
try:
existing = db.execute("SELECT amount FROM orders WHERE id=?", (order_id,)).fetchone()
if existing:
if existing[0] != amount:
raise ValueError("same identity with different intent")
db.execute("COMMIT")
return "existing"
db.execute("INSERT INTO orders VALUES (?, ?)", (order_id, amount))
if abort:
raise RuntimeError("injected failure before outbox write")
db.execute("INSERT INTO outbox VALUES (?, ?, 0)", (order_id, amount))
db.execute("COMMIT")
return "created"
except Exception:
db.execute("ROLLBACK")
raise
external_effects = {}
delivery_calls = []
def fictional_receiver(order_id, amount):
delivery_calls.append((order_id, amount))
if order_id in external_effects:
if external_effects[order_id] != amount:
raise ValueError("receiver identity collision")
return "existing"
external_effects[order_id] = amount
return "applied"
def relay(order_id, lose_response=False):
row = db.execute("SELECT amount, sent FROM outbox WHERE id=?", (order_id,)).fetchone()
if row is None:
raise LookupError("no publication intent")
if row[1]:
return "already-sent"
receipt = fictional_receiver(order_id, row[0])
if lose_response:
raise TimeoutError("injected response loss after receiver applied effect")
db.execute("UPDATE outbox SET sent=1 WHERE id=?", (order_id,))
return receipt
try:
check("unsafe attempt aborts", fails(lambda: unsafe_attempt("U1", 80, True), RuntimeError))
check("aborted database row absent", db.execute("SELECT * FROM orders WHERE id='U1'").fetchone() is None)
check("external effect survives rollback", unsafe_receipts == [("U1", 80)])
unsafe_attempt("U1", 80)
check("retry creates second external effect", len(unsafe_receipts) == 2)
check("one local row does not prove one external effect", db.execute("SELECT count(*) FROM orders WHERE id='U1'").fetchone()[0] == 1)
check("safe enqueue abort injected", fails(lambda: enqueue("S1", 120, True), RuntimeError))
check("order rolls back with failed intent", db.execute("SELECT * FROM orders WHERE id='S1'").fetchone() is None)
check("no partial outbox row", db.execute("SELECT * FROM outbox WHERE id='S1'").fetchone() is None)
check("successful atomic enqueue", enqueue("S1", 120) == "created")
check("order and intent both present", db.execute("SELECT o.amount, x.amount FROM orders o JOIN outbox x USING(id) WHERE id='S1'").fetchone() == (120, 120))
check("repeat original identity returns existing", enqueue("S1", 120) == "existing")
check("different amount under same identity rejected", fails(lambda: enqueue("S1", 121), ValueError))
check("collision does not change intent", db.execute("SELECT amount FROM outbox WHERE id='S1'").fetchone()[0] == 120)
check("response loss follows external effect", fails(lambda: relay("S1", True), TimeoutError))
check("effect exists despite timeout", external_effects == {"S1": 120})
check("intent remains retryable", db.execute("SELECT sent FROM outbox WHERE id='S1'").fetchone()[0] == 0)
check("receiver recognizes retry", relay("S1") == "existing")
check("two calls one effect", len(delivery_calls) == 2 and len(external_effects) == 1)
check("successful response marks sent", db.execute("SELECT sent FROM outbox WHERE id='S1'").fetchone()[0] == 1)
check("sent intent is not sent again", relay("S1") == "already-sent" and len(delivery_calls) == 2)
check("new ID creates distinct intent", enqueue("S2", 120) == "created")
check("new ID is a new external effect", relay("S2") == "applied" and len(external_effects) == 2)
check("receiver rejects mismatched retry payload", fails(lambda: fictional_receiver("S1", 999), ValueError))
check("receiver preserved original effect", external_effects["S1"] == 120)
check("relay rejects missing intent", fails(lambda: relay("missing"), LookupError))
print(json.dumps({"passed": len(checks), "checks": checks,
"python": platform.python_version(), "sqlite": sqlite3.sqlite_version,
"database": ":memory:", "network": False,
"vendorExecution": False,
"scope": "Sequential SQLite atomic writes and synthetic receiver identity checks; no concurrency, actual transport, Spanner isolation or durable external idempotency tested."}))
finally:
db.close()
Uma ordem foi confirmada uma vez na base, mas o callback repetido produziu dois recibos no processador externo. A reconciliação precisa de ambos os sistemas.
Armadilhas comuns
Equiparar timeout a rollback, acknowledgment a efeito único, nome de objeto a versão e FinalizeWriteStream a publicação do lote.
Tópicos relacionados: Migração e reconciliação de dados · Idempotência e recuperação de incidentes · Critérios de aceitação e handover para RUN
Identifica a fronteira de cada garantia, conserva identidade estável e exige evidência de conclusão em todos os sistemas relevantes.
Referência: Spanner transactions overview · Current linked standard guide; edition date unconfirmed (2026-09-30 inspection)