Data Platforms
Plataforma de dados fim a fim: captura mudanças do Postgres por replicação lógica, materializa camadas bruta, tratada e de consumo em colunar particionado e entrega o resultado por SQL, painéis e APIs com linhagem de coluna.
Introdução
Este guia descreve o Anthares Data Platform (a44-dp), o conjunto de serviços que leva um registro do instante em que é gravado numa tabela Postgres até o instante em que aparece num painel, numa resposta de API ou num arquivo exportado. Cobre captura de mudança, materialização em camadas, versionamento de modelos, testes, retenção, restauração pontual e autorização de leitura — cada estágio com parâmetros próprios e comportamento definido em falha.
A premissa é separar o sistema que atende usuários do que responde perguntas. O Postgres da aplicação segue otimizado para escritas curtas e índices B-tree; a plataforma nunca faz varredura analítica contra ele. Toda leitura pesada ocorre sobre cópias colunares alimentadas por change data capture. O custo é latência de 40 segundos a 15 minutos até o consumo; o ganho é que uma consulta sobre 400 milhões de linhas não degrada o checkout.
Aqui tratamos só do caminho do dado analítico. Regras reativas, filas de retentativa e orquestração estão em Automation — se você precisa reagir a um evento em menos de um segundo, é aquela página.
Os exemplos usam a CLI a44 3.8, o projeto acme e a região sa-east-1. Comandos Beta podem mudar de assinatura sem os 90 dias usuais de depreciação.
Visão geral
São cinco serviços com contratos explícitos entre si. O a44-capture mantém slots lógicos na origem e converte WAL em registros de mudança. O a44-land grava esses registros na bruta sem transformação, só anexando metadados. O a44-build roda os modelos versionados que produzem as camadas tratada e de consumo. O a44-query lê o colunar e responde a painéis, notebooks e APIs. O a44-catalog guarda esquemas, linhagem e políticas.
Cada camada tem contrato de imutabilidade, retenção e público distintos. A bruta é apenas-anexação e nunca corrigida — dado errado na origem permanece errado ali, por ser o registro histórico do que a origem informou. A tratada aplica tipagem, deduplicação por chave de negócio e normalização para UTC. A de consumo agrega, denormaliza e é a única visível externamente, o que permite reescrever modelos intermediários sem quebrar painéis.
_op, _lsn e _captured_at. Retenção de 90 dias em disco quente.Captura incremental
Replicação lógica com pgoutput, sem consulta periódica nem coluna de controle na origem.
Colunar particionado
Parquet com ZSTD nível 3, particionado por data de evento, grupos de linha de 64 MiB.
Linhagem de coluna
Cada coluna de consumo aponta para as de origem que a produziram, extraídas do SQL do modelo.
Quem já tem um lake externo pode provisionar só o a44-capture e escrever em bucket próprio; aí a linhagem cobre apenas até a fronteira do serviço e as tabelas externas viram nós opacos.
Primeiros passos
A configuração mínima exige um Postgres de origem com wal_level=logical, um destino de armazenamento e ao menos um modelo. O comando a44 data init gera a estrutura esperada pelo a44-build: models/, tests/ e platform.yaml. A plataforma não descobre tabelas sozinha — você declara quais entram na publicação, porque incluir tudo produz WAL desnecessário e expõe colunas que não deveriam sair do OLTP.
O primeiro snapshot copia as tabelas publicadas com REPEATABLE READ e instantâneo exportado alinhado ao LSN inicial do slot. Acima de 50 GB a cópia é fatiada por faixa de chave primária em 4 conexões. O slot já acumula WAL durante a cópia; se ela demorar mais que o espaço livre permite, a captura entra em slot_lag_critical e pausa para não encher o disco da origem.
Habilitar replicação lógica
wal_level=logical,max_replication_slotscom folga de 2 slots emax_wal_sendersequivalente. Mudarwal_levelexige reinício: 20 a 40 segundos em instâncias gerenciadas.Criar publicação e usuário de captura
Papel dedicado com
REPLICATIONeSELECTnas tabelas publicadas. Não reutilize o usuário da aplicação: revogado numa rotação de credencial, a captura para em silêncio.Registrar a origem
a44 data source addvalida conectividade, versão (mínimo Postgres 14) e tabelas declaradas. Falha na hora se alguma não tiver chave primária ouREPLICA IDENTITY FULL.Disparar o snapshot inicial
O progresso sai por tabela em linhas copiadas e bytes lidos. Um conjunto de 200 GB leva de 45 a 90 minutos com concorrência 4.
Construir o primeiro modelo
a44 data build --select tag:consumomaterializa só os modelos marcados, validando o caminho completo antes do cronograma automático.
a44 data init --project acme --region sa-east-1
a44 data source add postgres-main \
--dsn "postgresql://a44_capture@db-main.acme.internal:5432/app" \
--publication a44_pub \
--tables public.pedidos,public.itens_pedido,public.clientes \
--replica-identity-check strict
a44 data snapshot start postgres-main --parallel 4 --chunk-size 2GB
a44 data snapshot status postgres-main --watch
a44 data build --select tag:consumo --full-refresh
a44 data lineage show consumo.fato_pedidos.receita_liquida
Slot não consumido faz o Postgres reter WAL indefinidamente. Removendo a origem pela CLI, o slot cai junto; por outro caminho, execute SELECT pg_drop_replication_slot('a44_acme_main') manualmente. Defina max_slot_wal_keep_size em 20 GB: perder a captura e refazer o snapshot é preferível a derrubar o banco transacional.
Arquitetura
O dado atravessa quatro fronteiras de durabilidade: o WAL, onde a mudança já está confirmada antes de a plataforma vê-la; o buffer de captura, log em disco de 8 GiB por origem que absorve indisponibilidade do armazenamento sem replay do slot; o arquivo bruto no object storage, ponto em que o LSN é confirmado e o WAL liberado; e a materialização do consumo. Falha após a terceira fronteira nunca exige reler o banco transacional.
A confirmação de LSN é deliberadamente atrasada: o a44-capture só avança o confirmed_flush_lsn após o arquivo bruto ser gravado e o storage confirmar durabilidade, com no mínimo 10 segundos entre avanços. A origem retém, no pior caso, um lote mais 10 segundos — cerca de 2,4 GB a 30 MB/s. Reduzir o intervalo de lote corta latência, mas multiplica arquivos pequenos e piora a leitura colunar.
| Componente | Estado que mantém | Se cair sozinho | Tempo de recuperação |
|---|---|---|---|
a44-capture | Posição LSN e buffer local | WAL acumula na origem; nada se perde | ~15 s, do último LSN confirmado |
a44-land | Nenhum | Capture segura no buffer até 8 GiB | Imediato; qualquer réplica assume |
a44-build | Grafo e marcas d'água | Camadas param; bruta continua | 3 a 20 min, reexecuta o lote |
a44-query | Cache de metadados e resultado | Painéis e APIs retornam 503 | ~40 s, cache refeito sob demanda |
a44-catalog | Esquemas, linhagem, políticas | Builds não iniciam; leituras seguem 5 min com política em cache | ~60 s, de réplica síncrona |
O a44-catalog usa instância Postgres própria, com réplica síncrona na mesma região e assíncrona em outra — o que define o RPO discutido em Backup.
Banco de dados
A origem exige configuração explícita para que a captura seja correta, não apenas funcional. O ponto mais ignorado é REPLICA IDENTITY: no padrão DEFAULT, um UPDATE emite no WAL só a chave primária e as colunas alteradas. Isso basta para reconstruir o estado atual, mas impede detectar transições — saber que status saiu de pendente para pago exige REPLICA IDENTITY FULL, que aumenta o WAL entre 2x e 5x conforme a largura da linha.
DELETE e TRUNCATE recebem tratamento distinto. O DELETE vira uma linha com _op = 'D' na bruta, preservando o histórico, e a tratada marca _deleted_at. O TRUNCATE é decodificado como evento de tabela inteira e por padrão a plataforma o rejeita, parando a captura daquela tabela, porque quase sempre indica migração não coordenada. Para permitir, use on_truncate: propagate.
-- papel dedicado, sem acesso de escrita
CREATE ROLE a44_capture WITH LOGIN REPLICATION PASSWORD :'senha';
GRANT USAGE ON SCHEMA public TO a44_capture;
GRANT SELECT ON public.pedidos, public.itens_pedido, public.clientes TO a44_capture;
-- publicacao restrita: nunca FOR ALL TABLES em producao
CREATE PUBLICATION a44_pub
FOR TABLE public.pedidos, public.itens_pedido, public.clientes;
-- necessario para capturar valores antigos em UPDATE/DELETE
ALTER TABLE public.pedidos REPLICA IDENTITY FULL;
-- limite de retencao: invalida o slot antes de encher o volume
ALTER SYSTEM SET max_slot_wal_keep_size = '20GB';
-- diagnostico: atraso do slot em bytes
SELECT slot_name,
active,
pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn)) AS atraso
FROM pg_replication_slots
WHERE slot_name LIKE 'a44_%';
Replicação lógica não propaga DDL. Ao adicionar coluna, o decodificador emite o formato antigo até a plataforma detectar a divergência no lote seguinte — 30 a 90 segundos em que a coluna nova chega nula. DROP COLUMN ou tipo incompatível colocam a tabela em schema_drift e param só a captura dela. A retomada exige a44 data schema accept, que registra a versão no catálogo e reescreve o mapeamento de colunas.
Failover da origem é o outro ponto delicado: antes do Postgres 17 o slot não é replicado e uma promoção o destrói, obrigando a novo snapshot. Da versão 17 em diante o slot é criado com failover = true e a captura retoma em até 90 segundos.
Ingestão
Além do CDC há duas entradas: arquivos em prefixo monitorado e envio por endpoint HTTP. As três terminam na bruta com os mesmos metadados obrigatórios, mas diferem em ordem. O CDC preserva ordem transacional por tabela, derivada do LSN. Arquivos seguem a chegada do evento de criação, sem garantia entre arquivos concorrentes. O endpoint HTTP não garante ordem alguma e exige campo de sequência do produtor quando ela importa.
A deduplicação ocorre na tratada, não na ingestão. Cada registro carrega _ingest_id, um UUIDv7 gerado na entrada, e a chave natural do modelo. Se o mesmo evento chegar duas vezes — comum após retomada de captura, cuja garantia é ao menos uma vez — a bruta terá duas linhas e a tratada manterá só a de maior _lsn. Consultar a bruta pode, portanto, contar duplicados: é o comportamento pretendido e a razão de ela não ser exposta.
version: 3
sources:
postgres-main:
type: postgres_cdc
publication: a44_pub
slot: a44_acme_main
on_truncate: reject # reject | propagate | ignore
batch:
max_bytes: 256MiB
max_interval: 45s # menor valor aceito: 10s
max_rows: 500000
schema_drift: halt_table # halt_table | halt_source | auto_accept
arquivos-parceiro:
type: object_drop
prefix: s3://acme-inbox/parceiro/
format: csv
delimiter: ";"
encoding: latin1
on_bad_row: quarantine # quarantine | fail_batch | skip
max_bad_row_ratio: 0.02 # acima disso o lote inteiro falha
landing:
destination: s3://acme-lake/bruta/
file_format: parquet
compression: zstd:3
target_file_size: 128MiB
partition_by: [_event_date, _source]
from anthares44 import DataPlatform
dp = DataPlatform(project="acme", token=os.environ["A44_TOKEN"])
stream = dp.stream("eventos_web")
# lote de ate 5.000 registros ou 4 MiB, o que vier primeiro
with stream.batch(max_rows=5000, max_bytes=4 * 1024 * 1024) as lote:
for ev in ler_eventos():
lote.append({
"evento_id": ev.id, # chave natural, usada na deduplicacao
"seq": ev.sequencia, # define ordem quando o transporte nao define
"ocorrido_em": ev.ts.isoformat(),
"payload": ev.dados,
})
# resposta: aceitos, rejeitados por esquema e id do lote na camada bruta
print(lote.result.accepted, lote.result.rejected, lote.result.batch_id)
# rejeicoes ficam disponiveis por 14 dias com o motivo por linha
for r in lote.result.rejections(limit=20):
print(r.line, r.reason) # ex.: 3, "campo 'ocorrido_em' nao e ISO-8601"
{
"_ingest_id": "01912f3a-7b41-7c9e-9d20-5a1f0b7c4e88",
"_source": "postgres-main",
"_table": "public.pedidos",
"_op": "U",
"_lsn": "3A/7F0128C0",
"_tx_id": 884210377,
"_captured_at": "2026-07-19T14:02:11.418Z",
"_landed_at": "2026-07-19T14:02:53.902Z",
"_schema_version": 7,
"before": { "status": "pendente", "valor_total": "310.00" },
"after": { "id": 99120, "status": "pago", "valor_total": "310.00",
"cliente_id": 4471, "atualizado_em": "2026-07-19T14:02:11Z" }
}
Linhas com o mesmo _tx_id foram confirmadas juntas na origem. Modelos que exigem consistência entre tabelas devem filtrar pela marca d'água _watermark_lsn, nunca pelo timestamp de captura, que não é monotônico entre tabelas.
Analytics
O a44-query é um motor vetorizado que lê Parquet direto do object storage, mantendo em memória apenas metadados de arquivo e resultados intermediários. Não há carga prévia: a tabela é consultável assim que o manifesto é publicado. O planejador poda em três etapas nesta ordem — eliminação de partição pelo predicado de data, descarte de grupo de linha por mínimo e máximo, e projeção de coluna. Filtrar sete dias numa tabela de dois anos lendo 3 de 40 colunas costuma tocar menos de 1% dos bytes.
A ordem física dentro da partição importa tanto quanto o particionamento: como o descarte depende de mínimo e máximo, colunas dispersas aleatoriamente não podam nada. Daí o cluster_by com até três colunas, que reordena os dados na escrita: custa de 20% a 60% a mais de materialização e devolve de 4x a 15x menos bytes lidos nas consultas filtradas por elas. Compensa com filtro estável, não em exploração livre.
-- particionada por data_evento, agrupada por (loja_id, categoria)
EXPLAIN (ANALYZE, FORMAT TEXT)
SELECT loja_id,
categoria,
date_trunc('day', data_evento) AS dia,
count(*) AS pedidos,
sum(receita_liquida) AS receita,
approx_percentile(ticket, 0.95) AS p95_ticket
FROM consumo.fato_pedidos
WHERE data_evento >= current_date - INTERVAL '7 days'
AND loja_id IN (11, 42, 77)
GROUP BY 1, 2, 3
ORDER BY receita DESC
LIMIT 200;
-- saida resumida:
-- particoes lidas: 7 de 731
-- grupos de linha: 42 de 5.118 (descarte por min/max em loja_id)
-- bytes varridos: 318 MiB de 214 GiB
-- tempo: 1.24 s memoria de pico: 640 MiB
| Parâmetro | Padrão | Descrição |
|---|---|---|
query_timeout | 300s | Devolve QUERY_TIMEOUT. Máximo de 3600s em pools dedicados. |
max_scan_bytes | 2TiB | Aborta antes de executar se a estimativa for maior. Protege contra WHERE esquecido. |
spill_to_disk | true | Derrama agregações e junções para disco. Fica 3x a 8x mais lento, mas evita OOM. |
result_cache_ttl | 900s | Invalidado quando o manifesto de qualquer tabela referenciada muda. |
max_concurrent_queries | 16 | Por pool. Acima de 64 na fila, novas requisições recebem 429. |
approx_threshold | 10000000 | Acima disso, count(distinct) usa HyperLogLog com erro de ~1,6%. |
Em junções, o lado menor vira tabela de dispersão em memória; se não couber no pool, o planejador troca para junção particionada, que exige duas passagens e costuma triplicar o tempo. a44 data query profile mostra a estratégia escolhida e os bytes derramados.
Storage
Todo dado persistente vive em object storage compatível com S3, com um prefixo por camada e um manifesto por tabela. O manifesto é um JSON versionado que lista os arquivos válidos daquela versão, com estatísticas por coluna. Publicar uma versão é uma escrita atômica de manifesto; leitores que já abriram a anterior seguem nos arquivos antigos até terminar. É isso que permite reescrever uma partição inteira sem que um painel veja resultado parcial.
Arquivos órfãos — fora de qualquer manifesto retido — são removidos por varredura diária com carência de 72 horas, já que uma consulta longa pode ainda estar lendo um arquivo recém-descartado. Quem reescreve tabelas grandes com frequência ocupa, por até três dias, mais espaço que a soma das versões vivas.
| Camada | Classe de armazenamento | Retenção | Versões de manifesto | Após a retenção |
|---|---|---|---|---|
| Bruta | Padrão (quente) | 90 dias | 30 | Vai para acesso infrequente; leitura passa a custar por GB |
| Bruta arquivada | Arquivo frio | 7 anos | 1 | Expurgo definitivo; reidratação de 3 a 12 h |
| Tratada | Padrão | 24 meses | 60 | Partições apagadas, recomputáveis da bruta arquivada |
| Consumo | Padrão + cache SSD | 5 anos | 90 | Agregados mensais mantidos, grão diário descartado |
| Resultados de consulta | Efêmero | 15 dias | — | Apagado sem aviso; não use como armazenamento |
# quanto cada camada ocupa, com e sem versoes historicas
a44 data storage usage --by-layer --include-orphans
# LAYER ATIVO VERSOES ORFAOS TOTAL
# bruta 4.1 TiB 612 GiB 88 GiB 4.8 TiB
# tratada 1.7 TiB 240 GiB 12 GiB 1.9 TiB
# consumo 390 GiB 77 GiB 2 GiB 469 GiB
# compactar arquivos pequenos: alvo de 128 MiB por arquivo
a44 data compact consumo.fato_pedidos \
--partition "data_evento >= '2026-06-01'" \
--target-file-size 128MiB \
--max-concurrency 8
# antecipar limpeza de orfaos (respeita carencia de 72h)
a44 data vacuum --layer bruta --dry-run
Cada manifesto carrega checksum SHA-256 do conjunto de arquivos que descreve. O a44-query o valida ao abrir a tabela e recusa manifestos inconsistentes com MANIFEST_CHECKSUM_MISMATCH, em vez de devolver resultado parcial em silêncio.
Data lake
O particionamento padrão é por data de evento, não de ingestão. A diferença aparece em correções tardias: um pedido de 3 de junho corrigido em 19 de julho reescreve a partição de 3 de junho. Consultas por período de negócio ficam corretas sem filtro extra, mas partições antigas deixam de ser imutáveis. Para imutabilidade histórica — conciliação contábil fechada — use particionamento duplo por data de evento e de processamento.
O tamanho de arquivo é o principal regulador de desempenho do lake: arquivos pequenos multiplicam listagens e leituras de rodapé, grandes demais reduzem paralelismo e encarecem a reescrita da partição. A faixa recomendada vai de 64 MiB a 256 MiB, padrão 128 MiB, grupos de linha de 64 MiB. O aviso SMALL_FILES aparece quando mais de 30% dos arquivos ficam abaixo de 16 MiB. Uma partição com 20.000 arquivos de 2 MiB leva 40 segundos só para planejar, contra menos de 1 segundo com 300 de 128 MiB.
_event_date e chaves adicionais. Unidade de reescrita e de expurgo.from a44_build import model, ref, test
@model(
layer="consumo",
materialization="incremental",
unique_key=["pedido_id"],
partition_by="data_evento",
cluster_by=["loja_id", "categoria"],
file_format="parquet",
target_file_size="128MiB",
# reprocessa 3 dias para absorver correcoes tardias da origem
lookback_window="3 days",
tags=["consumo", "critico"],
)
def fato_pedidos(ctx):
pedidos = ref("tratada.pedidos")
itens = ref("tratada.itens_pedido")
if ctx.is_incremental:
pedidos = pedidos.filter(
ctx.col("atualizado_em") >= ctx.watermark - ctx.lookback
)
return (
pedidos.join(itens, on="pedido_id", how="left")
.group_by("pedido_id", "loja_id", "categoria", "data_evento")
.agg(
receita_liquida=ctx.sum("valor_item") - ctx.sum("desconto"),
ticket=ctx.sum("valor_item"),
itens=ctx.count("item_id"),
)
)
test(fato_pedidos, "unicidade", columns=["pedido_id"])
test(fato_pedidos, "frescor", column="data_evento", max_atraso="90 minutes")
test(fato_pedidos, "faixa", column="receita_liquida", min=-50000, max=2000000)
lookback_window define quantos dias de partição são reescritos por execução incremental. Três dias numa tabela de 200 GB reescrevem cerca de 800 MB por ciclo; 30 dias multiplicam isso por dez e, a cada 15 minutos, saem mais caro que uma reconstrução completa diária. Janela longa exige cadência maior.
Dashboards
Painéis são arquivos versionados no repositório dos modelos, o que impede um gráfico de referenciar coluna inexistente: o a44-build valida as referências contra o catálogo e falha a publicação com DASHBOARD_BROKEN_REF. Cada visualização declara a tabela de consumo que lê, os filtros e a política de atualização. Referência à bruta ou à tratada é recusada na validação.
São três modos de atualização, e a escolha define custo e latência percebida. Em on_view a consulta roda a cada abertura, respeitando o result_cache_ttl — adequado a painéis com poucos acessos diários. Em scheduled ela roda em cronograma fixo e o painel serve o último resultado gravado, com leitura abaixo de 200 ms, o que atende muitos espectadores simultâneos. Em on_build a nova versão da tabela dispara a atualização, dando o menor atraso ao custo de execuções irregulares.
Marca de frescor
Exibe o instante do dado mais recente da tabela, não o da última execução da consulta.
Estado degradado
Teste de qualidade reprovado gera faixa de alerta no painel, em vez de esconder o problema.
Filtro obrigatório
Exigir filtro de período antes de executar evita varredura acidental de anos inteiros.
{
"slug": "receita-diaria",
"titulo": "Receita diária por loja",
"refresh": { "mode": "scheduled", "cron": "*/15 * * * *", "timezone": "America/Sao_Paulo" },
"required_filters": ["periodo"],
"filters": {
"periodo": { "type": "date_range", "default": "last_30_days", "max_span_days": 400 },
"loja_id": { "type": "multi_select", "source": "consumo.dim_loja.loja_id", "max_selected": 50 }
},
"tiles": [
{
"id": "serie_receita",
"type": "line",
"table": "consumo.fato_pedidos",
"x": "data_evento",
"y": { "column": "receita_liquida", "agg": "sum" },
"series": "loja_id",
"freshness_badge": true,
"degrade_on_test_failure": ["frescor", "unicidade"]
}
],
"access": { "grupos": ["financeiro", "operacoes"], "masking_profile": "padrao" }
}
max_span_days é o principal controle de custo em painéis abertos: sem ele um espectador seleciona cinco anos e dispara varredura de vários terabytes. A seleção é recusada com FILTER_SPAN_EXCEEDED.
APIs
Tabelas de consumo podem virar endpoints HTTP somente-leitura sem código de servidor. Cada publicação define colunas visíveis, filtros aceitos, limite de linhas por resposta e perfil de mascaramento. O endpoint fica em https://data.acme.a44.app/v1/<publicacao> e pagina por cursor opaco — nunca por deslocamento, que em tabela colunar grande obriga a varrer tudo até a posição pedida.
A autorização é por consumidor, não por usuário final. Cada chave pertence a um consumidor registrado que carrega colunas permitidas, perfil de mascaramento e cota. Dois consumidores da mesma publicação podem ver a mesma linha com valores diferentes: um recebe cpf completo, outro ***.***.789-01. A decisão é aplicada no motor, na projeção — não há como obter a coluna em claro pedindo de outra forma.
import { DataClient } from "@anthares44/data";
const client = new DataClient({
baseUrl: "https://data.acme.a44.app/v1",
apiKey: process.env.A44_DATA_KEY!,
timeoutMs: 30_000,
});
// paginacao por cursor: o cursor codifica a posicao fisica, nao um offset
let cursor: string | undefined;
let total = 0;
do {
const page = await client.query("pedidos_por_loja", {
filtros: { data_evento: { gte: "2026-07-01", lt: "2026-08-01" }, loja_id: [11, 42] },
colunas: ["loja_id", "data_evento", "receita_liquida"],
ordenar: [{ coluna: "data_evento", direcao: "asc" }],
limite: 5000,
cursor,
});
total += page.rows.length;
cursor = page.nextCursor; // undefined encerra a paginacao
// page.meta.scannedBytes e page.meta.freshnessAt vem em todas as respostas
} while (cursor);
console.log(total);
from anthares44.data import DataClient, RateLimited
import time
client = DataClient(base_url="https://data.acme.a44.app/v1",
api_key=os.environ["A44_DATA_KEY"],
timeout=30)
def paginar(publicacao, **kwargs):
cursor, tentativa = None, 0
while True:
try:
pagina = client.query(publicacao, cursor=cursor, limite=5000, **kwargs)
except RateLimited as e:
# backoff exponencial respeitando Retry-After do servidor
espera = e.retry_after or min(2 ** tentativa, 60)
time.sleep(espera)
tentativa += 1
continue
tentativa = 0
yield from pagina.rows
cursor = pagina.next_cursor
if cursor is None:
return
for linha in paginar("pedidos_por_loja",
filtros={"data_evento": {"gte": "2026-07-01"}},
colunas=["loja_id", "receita_liquida"]):
processar(linha)
| Código | Situação | Cabeçalho relevante | Ação recomendada |
|---|---|---|---|
200 | Página retornada | X-A44-Scanned-Bytes | Seguir com next_cursor até nulo |
400 | Filtro em coluna não publicada | — | Corrigir; retentar não resolve |
403 | Coluna fora do escopo | X-A44-Denied-Columns | Ajustar a projeção ou ampliar o escopo |
410 | Cursor expirado (10 min) | — | Reiniciar a paginação |
413 | Varredura acima da cota | X-A44-Estimated-Bytes | Reduzir datas ou colunas |
429 | Cota esgotada | Retry-After | Backoff exponencial, teto de 60 s |
As cotas por consumidor têm três dimensões independentes: 600 requisições por minuto, 200 GiB varridos por hora e 8 consultas simultâneas. Estourar qualquer uma produz 429 com X-A44-Quota-Dimension indicando qual.
Backup
São três objetos de recuperação com garantias distintas. A origem transacional tem backup contínuo, RPO de 5 minutos e restauração pontual em qualquer segundo dentro de 35 dias. O lake tem versionamento de objeto e replicação para uma segunda região, com RPO de 15 minutos. O catálogo tem réplica síncrona na região primária e RPO efetivo de zero para falha de instância única.
Restauração pontual do lake não copia arquivos: reconstrói o estado escolhendo o manifesto vigente no instante alvo, o que leva segundos dentro da janela de versões retidas — 90 versões no consumo, cerca de 22 horas a 15 minutos de cadência. Fora dela é preciso recomputar os modelos da bruta, de 20 minutos a 6 horas, ou reidratar a bruta arquivada, somando de 3 a 12 horas.
Identificar o instante alvo
a44 data historylista versões de manifesto com data, arquivos, linhas e o build que as produziu. Escolha a última anterior ao incidente.Validar antes de aplicar
Consulte a versão antiga por viagem no tempo, sem restaurar. Evita restaurar sobre um estado que já estava errado.
Represar o cronograma
Pause o agendamento antes de restaurar; se o build rodar no meio, sobrescreve a restauração com dados da mesma origem defeituosa.
Restaurar publicando manifesto
A restauração cria uma versão nova idêntica à antiga em vez de apagar as posteriores, mantendo o incidente auditável.
Registrar o RPO observado
Compare o instante do último dado íntegro com o do incidente. Esse número, e não o alvo contratual, alimenta a revisão pós-incidente.
a44 data schedule pause --select tag:consumo
a44 data history consumo.fato_pedidos --limit 5
# VERSAO PUBLICADO_EM ARQUIVOS LINHAS BUILD
# 4182 2026-07-19T14:15:02Z 5.118 412.907.331 bld_9f2a1c
# 4181 2026-07-19T14:00:04Z 5.114 412.884.020 bld_9f2a0b
# ler a versao anterior sem restaurar
a44 data query "select count(*), sum(receita_liquida)
from consumo.fato_pedidos for version as of 4181"
a44 data restore consumo.fato_pedidos --to-version 4181 --reason "INC-2291"
a44 data schedule resume --select tag:consumo
Um agendamento embutido restaura mensalmente uma tabela tag:critico em ambiente isolado, compara contagens e checksums e registra o tempo total. Backup nunca restaurado é hipótese, não garantia; o relatório sai em a44 data restore-drill report.
Escalabilidade
Os limites não são um número único. A captura satura por taxa de WAL: um a44-capture sustenta cerca de 180 MB/s decodificados, algo entre 40.000 e 90.000 mudanças por segundo conforme a largura das linhas. Acima disso o atraso do slot cresce sem convergir e a saída é dividir a publicação em slots independentes, o que duplica a leitura do WAL e soma 10% a 15% de CPU na origem.
A materialização satura por volume reescrito por ciclo, não por tamanho da tabela. Um build incremental de até 40 GB termina em menos de 10 minutos com 16 trabalhadores. Acima de 150 GB por ciclo a escrita excede a cadência de 15 minutos, as execuções se sobrepõem e a plataforma emite BUILD_OVERLAP em vez de iniciar o ciclo seguinte. A correção quase nunca é adicionar trabalhadores: é reduzir lookback_window, particionar mais fino ou alongar a cadência.
| Dimensão | Limite confortável | Teto observado | O que quebra primeiro depois disso |
|---|---|---|---|
| Taxa de WAL por slot | 120 MB/s | 180 MB/s | Atraso cresce sem convergir; risco de invalidação |
| Tabelas por publicação | 150 | 400 | Recarga de esquema após DDL passa de 90 s e trava o lote |
| Volume reescrito por ciclo | 40 GB | 150 GB | BUILD_OVERLAP; a fila cresce |
| Partições por tabela | 20.000 | 120.000 | Planejamento passa de 5 s só para listar |
| Modelos no grafo | 800 | 2.500 | Resolver dependências domina o build |
| Consultas por pool | 16 | 64 | Derramamento geral; latência 8x pior |
| Consumidores de API | 200 | 1.000 | Avaliar política por requisição vira gargalo de CPU |
Réplicas de a44-capture não aumentam a vazão de um slot: a decodificação lógica é sequencial por natureza, porque a ordem transacional depende dela. Réplicas só assumem em caso de falha. Se a origem gera mais WAL do que um slot decodifica, resta dividir a publicação — e isso quebra a garantia de ordem entre as tabelas separadas.
Monitoramento
Quatro métricas explicam a maior parte dos incidentes e devem ser instrumentadas antes das demais. capture_slot_lag_bytes mede o WAL retido na origem e é o único sinal capaz de derrubar o banco transacional se ignorado. freshness_seconds mede a distância entre agora e o dado mais recente materializado — o que o usuário percebe. build_duration_seconds contra a cadência antecipa sobreposição de ciclos. query_scanned_bytes por consumidor mostra quem gasta a cota antes que a fatura explique.
Alertas devem ser definidos por sintoma, não por componente: "CPU alta no serviço" é ruído, "frescor da tabela crítica acima de 90 minutos" é algo a resolver. As métricas saem em OpenMetrics em /metrics e os eventos por webhook, com deduplicação por chave de incidente para evitar avalanche quando uma origem cai e cinquenta tabelas ficam desatualizadas juntas.
alertas:
- nome: slot_lag_critico
expr: capture_slot_lag_bytes > 15e9 # 15 GB, antes do teto de 20 GB
for: 5m
severidade: pagina
dedup_key: "{{ source }}"
runbook: https://docs.acme.internal/runbooks/slot-lag
- nome: frescor_degradado
expr: freshness_seconds{tag="critico"} > 5400 # 90 minutos
for: 10m
severidade: pagina
dedup_key: "frescor/{{ source }}" # agrupa tabelas da mesma origem
- nome: build_perto_da_cadencia
expr: build_duration_seconds / build_interval_seconds > 0.8
for: 3 ciclos
severidade: aviso
- nome: consumidor_estourando_cota
expr: rate(query_scanned_bytes[1h]) > 150e9
for: 15m
severidade: aviso
incluir_labels: [consumer_id, publicacao]
destinos:
pagina: { webhook: https://ops.acme.internal/a44, repetir_apos: 20m }
aviso: { webhook: https://ops.acme.internal/a44, repetir_apos: 6h }
Frescor por tabela
Medido contra o evento na origem, não contra o build. Detecta captura parada com builds verdes.
Orçamento de ciclo
Duração do build sobre cadência. Acima de 0,8 por três ciclos, a sobreposição é questão de tempo.
Custo atribuído
Bytes varridos e tempo de motor por consumidor, publicação e painel, para cobrança interna.
Eventos de webhook trazem incident_key, first_seen_at e affected_objects. Uma origem que cai gera um evento único com todas as tabelas afetadas, não um por tabela — o receptor precisa saber expandir a lista.
Qualidade
Testes são declarados junto do modelo e rodam após a materialização, antes de publicar o manifesto. Falha error impede a publicação: a tabela fica na versão anterior e os consumidores veem dados antigos, porém corretos. Com warn a publicação prossegue e o painel exibe a faixa degradada descrita em Dashboards. Calibrar a severidade por teste é o principal controle de risco: bloquear tudo para a plataforma por uma coluna secundária, bloquear nada torna os testes decorativos.
Quatro famílias cobrem a maioria dos casos. Frescor compara o máximo da coluna de tempo com o instante atual. Unicidade conta chaves duplicadas — em CDC, duplicidade quase sempre indica falha na resolução por _lsn. Faixa verifica limites numéricos, útil contra sinal invertido ou unidade errada. Integridade referencial confere se toda chave estrangeira existe na dimensão, tolerando dimensões atrasadas.
testes:
- modelo: consumo.fato_pedidos
checagens:
- tipo: frescor
coluna: data_evento
max_atraso: 90m
severidade: error
janela_silencio: ["sab 00:00-23:59", "dom 00:00-23:59"]
- tipo: unicidade
colunas: [pedido_id]
severidade: error
amostra_falhas: 20 # grava ate 20 chaves duplicadas no relatorio
- tipo: faixa
coluna: receita_liquida
min: -50000 # devolucoes produzem valor negativo legitimo
max: 2000000
severidade: warn
tolerancia_linhas: 0.001 # ate 0,1% fora da faixa nao falha
- tipo: integridade
coluna: loja_id
referencia: consumo.dim_loja.loja_id
tolerancia_linhas: 0.005 # dimensao pode chegar ate 1 ciclo atrasada
severidade: error
- tipo: volume
comparar_com: media_7d
desvio_max: 0.35 # +/-35% de variacao no numero de linhas
severidade: warn
Quando um teste falha, a44 data lineage impact lista as colunas, painéis e publicações de API que dependem da coluna afetada. A linhagem vem da análise sintática do SQL, não de anotação manual, então acompanha refatorações. Limitação: SELECT * e SQL dinâmico resolvem só até o nível de tabela, marcados como coarse.
O histórico de execução é retido por 400 dias com a versão de manifesto vigente, o resultado de cada checagem e as amostras de falha — suficiente para auditar se a tabela passava nos testes na data de um fechamento contábil.
Segurança
O mascaramento acontece em dois momentos, e confundi-los produz vazamento. O estático ocorre na transição bruta para tratada: a coluna vira hash com sal por projeto e o valor em claro só permanece na bruta, restrita a um papel de quebra de emergência com registro obrigatório. O dinâmico ocorre na leitura, por perfil de consumidor, sobre dados em claro. Use estático para o que ninguém deve ver após a ingestão; dinâmico para o que alguns papéis precisam ver.
As políticas ficam no catálogo e são avaliadas pelo motor durante a projeção. Uma política nega colunas, não linhas: sem acesso a cpf, o consumidor recebe as demais normalmente e um 403 só ao pedir a negada. Filtros de linha são separados, como predicados injetados no plano — úteis para limitar um parceiro às próprias lojas. Ambos correm antes da leitura, de modo que um WHERE criativo não infere valores negados por mensagem de erro ou tempo de resposta.
-- mascaramento estatico: aplicado na bruta -> tratada, irreversivel
CREATE MASKING POLICY cpf_estatico
ON tratada.clientes.cpf
USING hash_sha256(valor || projeto_salt())
RETENTION raw_visible_to = 'papel:quebra_emergencia';
-- mascaramento dinamico: aplicado na leitura, por perfil de consumidor
CREATE MASKING POLICY telefone_dinamico
ON consumo.dim_cliente.telefone
USING CASE
WHEN current_profile() IN ('suporte_n2', 'antifraude') THEN valor
WHEN current_profile() = 'parceiro' THEN NULL
ELSE mask_keep_last(valor, 4)
END;
-- filtro de linha por consumidor: parceiro so ve as proprias lojas
CREATE ROW POLICY parceiro_lojas
ON consumo.fato_pedidos
FOR CONSUMER GROUP 'parceiro'
USING loja_id IN (SELECT loja_id FROM meta.lojas_do_consumidor
WHERE consumer_id = current_consumer());
-- escopo de leitura de um consumidor de API
GRANT SELECT (loja_id, data_evento, receita_liquida)
ON consumo.fato_pedidos
TO CONSUMER 'parceiro-logistica'
WITH MASKING PROFILE 'parceiro', QUOTA scanned_bytes_per_hour = '50GiB';
Mascarar cpf e liberar count(*) GROUP BY cpf_hash ainda permite contar e cruzar indivíduos, já que o hash determinístico funciona como identificador estável. Em consultas expostas a terceiros, combine mascaramento com MIN GROUP SIZE 25, que descarta grupos menores ao custo de tornar os totais não somáveis.
Chaves com escopo
Vinculadas a um consumidor, validade máxima de 180 dias e rotação com sobreposição de 48 h.
Trilha de auditoria
Cada leitura registra consumidor, colunas, política e bytes. Retenção imutável de 400 dias.
Isolamento de rede
Origem e serviços falam por VPN privada; o storage só aceita o endpoint interno do projeto.
Boas práticas
Trate modelos como código de produção sujeito a revisão. Alteração em modelo de consumo deve rodar em validação sobre amostra determinística: a44 data build --env staging --sample 5% materializa uma fração estável por hash da chave, permitindo comparar contagens e somas sem processar o conjunto inteiro. Diferenças acima de 0,5% em métricas monetárias devem bloquear a fusão; a comparação sai de a44 data diff.
Evite três padrões que geram incidentes recorrentes. O primeiro é derivar métrica de negócio direto da bruta para "ganhar latência" — funciona até a primeira duplicata de retomada de captura. O segundo é SELECT * em modelos, que quebra a linhagem de coluna e propaga mudanças de esquema sem revisão. O terceiro é marcar todos os testes como error, o que leva a equipe a desligá-los em massa no primeiro incidente noturno em vez de calibrar severidades.
Uma definição por métrica
Receita líquida calculada em dois modelos diverge. Defina uma vez no consumo e referencie.
Cadência por criticidade
Poucas tabelas precisam de 15 minutos. Passar as secundárias para 6 horas libera ~40% do build.
Nomenclatura previsível
fato_, dim_ e agg_ com o grão no nome evitam junções em grãos incompatíveis.
Depreciação anunciada
Marque a coluna deprecated 60 dias antes de remover; o catálogo avisa quem a lê.
#!/usr/bin/env bash
set -euo pipefail
# 1. sintaxe, referencias e ciclos no grafo
a44 data compile --strict --fail-on-warning
# 2. modelos afetados pelo diff, nao o projeto inteiro
ALVOS=$(a44 data ls --changed-since origin/main --format selector)
[ -z "$ALVOS" ] && { echo "nenhum modelo alterado"; exit 0; }
# 3. materializa amostra deterministica em staging
a44 data build --env staging --select "$ALVOS+" --sample 5% --seed 44
# 4. compara metricas com producao; falha acima de 0,5%
a44 data diff --env staging --baseline prod \
--select "$ALVOS" \
--metrics "count(*),sum(receita_liquida),count(distinct pedido_id)" \
--max-drift 0.005
# 5. impacto a jusante: quem quebra se isto for publicado
a44 data lineage impact --select "$ALVOS" --format markdown > impacto.md
O sufixo + em --select inclui os dependentes a jusante; sem ele você valida o modelo alterado, mas não os agregados que o consomem — onde a maioria das regressões aparece.
FAQ
As perguntas abaixo complementam as seções anteriores em vez de resumi-las; quando a resposta depende de um parâmetro já descrito, o link aponta para sua definição.
Origens não-Postgres, residência fora de sa-east-1 e catálogos externos ficam fora deste guia: alteram as garantias de ordem e de RPO aqui descritas.
Quanto tempo leva para um dado gravado no Postgres aparecer em um painel?
Entre 40 segundos e 15 minutos: até 45 s no lote de captura (batch.max_interval), 5 a 20 s de escrita na bruta e o restante até o próximo ciclo de build. Baixar max_interval para 10 s corta a primeira parcela, mas produz arquivos menores e encarece a compactação.
Posso consultar a camada bruta diretamente?
Sim, com papel específico, mas os resultados podem trazer duplicatas de retomada e valores pré-mascaramento. Ela não é exposta a APIs nem a painéis por decisão de projeto. Para latência menor que o ciclo de build, crie um modelo de consumo com cadência curta.
O que acontece se a origem sofrer failover?
No Postgres 17 ou superior, com slot failover = true sincronizado antes do evento, a captura retoma em até 90 segundos sem perda. Em versões anteriores o slot morre na promoção e exige novo snapshot, com o tempo indicado em Primeiros passos.
Como reprocessar um período depois de corrigir um erro de modelo?
a44 data build --select modelo+ --backfill "2026-05-01..2026-06-30" reescreve só as partições do intervalo e propaga aos dependentes. Ele estima os bytes antes de executar e pede confirmação acima de 500 GB. Não altera a bruta, então é repetível.
Qual é o RPO real em uma perda total de região?
15 minutos para o lake, limitado pela replicação entre regiões, e 5 minutos para a origem transacional; o catálogo replica de forma assíncrona com atraso abaixo de 30 segundos. O restabelecimento completo em outra região leva de 40 a 90 minutos, dominado pela reconstrução do cache de metadados.
Por que meu build ficou mais lento sem eu ter mudado nada?
Quase sempre acúmulo de arquivos pequenos na partição corrente ou crescimento do volume dentro da lookback_window: verifique com a44 data storage usage --by-partition e compacte. A terceira causa é uma dimensão que deixou de caber em memória na junção, visível em a44 data query profile.
Dá para dar acesso a um parceiro externo sem expor o resto dos dados?
Sim: registre um consumidor, conceda as colunas necessárias com GRANT SELECT (colunas), aplique perfil de mascaramento e política de linha. As três camadas são avaliadas no motor, antes da leitura. Defina também cota de bytes por hora — sem ela um parceiro consome o orçamento do projeto em uma tarde.
A linhagem de coluna funciona com SQL gerado dinamicamente?
Parcialmente: a análise resolve referências explícitas, mas construções dinâmicas e SELECT * viram nós coarse, com linhagem só em nível de tabela — basta para impacto de remoção de tabela, não para rastrear uma coluna.