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.

Guia Atualizado em julho de 2026 Leitura: ~26 min

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.

Escopo deste documento

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.

ConsumoDenormalizada, por grão de negócio. Único nível exposto externamente. Retenção de 5 anos.
TratadaTipos canônicos, deduplicação, UTC, colunas sensíveis mascaradas. Retenção de 24 meses.
BrutaApenas-anexação, com _op, _lsn e _captured_at. Retenção de 90 dias em disco quente.
OrigemPostgres OLTP com slots lógicos e publicações por tabela. Nunca recebe carga analítica.

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.

  1. Habilitar replicação lógica

    wal_level=logical, max_replication_slots com folga de 2 slots e max_wal_senders equivalente. Mudar wal_level exige reinício: 20 a 40 segundos em instâncias gerenciadas.

  2. Criar publicação e usuário de captura

    Papel dedicado com REPLICATION e SELECT nas tabelas publicadas. Não reutilize o usuário da aplicação: revogado numa rotação de credencial, a captura para em silêncio.

  3. Registrar a origem

    a44 data source add valida conectividade, versão (mínimo Postgres 14) e tabelas declaradas. Falha na hora se alguma não tiver chave primária ou REPLICA IDENTITY FULL.

  4. 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.

  5. Construir o primeiro modelo

    a44 data build --select tag:consumo materializa só os modelos marcados, validando o caminho completo antes do cronograma automático.

terminal
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 órfão enche o disco da origem

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.

Postgres OLTPslot lógico
a44-capturebuffer 8 GiB
Camada brutaParquet append
a44-buildmodelos + testes
a44-querySQL / API
ComponenteEstado que mantémSe cair sozinhoTempo de recuperação
a44-capturePosição LSN e buffer localWAL acumula na origem; nada se perde~15 s, do último LSN confirmado
a44-landNenhumCapture segura no buffer até 8 GiBImediato; qualquer réplica assume
a44-buildGrafo e marcas d'águaCamadas param; bruta continua3 a 20 min, reexecuta o lote
a44-queryCache de metadados e resultadoPainéis e APIs retornam 503~40 s, cache refeito sob demanda
a44-catalogEsquemas, linhagem, políticasBuilds 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.

sql/preparar-origem.sql
-- 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_%';
Mudança de esquema na origem

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.

platform.yaml
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]
ingest_eventos.py
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"
registro na camada bruta (exemplo)
{
  "_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.

sql/consulta-com-poda.sql
-- 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âmetroPadrãoDescrição
query_timeout300sDevolve QUERY_TIMEOUT. Máximo de 3600s em pools dedicados.
max_scan_bytes2TiBAborta antes de executar se a estimativa for maior. Protege contra WHERE esquecido.
spill_to_disktrueDerrama agregações e junções para disco. Fica 3x a 8x mais lento, mas evita OOM.
result_cache_ttl900sInvalidado quando o manifesto de qualquer tabela referenciada muda.
max_concurrent_queries16Por pool. Acima de 64 na fila, novas requisições recebem 429.
approx_threshold10000000Acima 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.

CamadaClasse de armazenamentoRetençãoVersões de manifestoApós a retenção
BrutaPadrão (quente)90 dias30Vai para acesso infrequente; leitura passa a custar por GB
Bruta arquivadaArquivo frio7 anos1Expurgo definitivo; reidratação de 3 a 12 h
TratadaPadrão24 meses60Partições apagadas, recomputáveis da bruta arquivada
ConsumoPadrão + cache SSD5 anos90Agregados mensais mantidos, grão diário descartado
Resultados de consultaEfêmero15 diasApagado sem aviso; não use como armazenamento
terminal
# 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
Escrita atômica verificável

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.

ManifestoVersão, arquivos válidos, estatísticas por coluna, checksum. Escrita atômica.
PartiçãoPor _event_date e chaves adicionais. Unidade de reescrita e de expurgo.
Arquivo ParquetAlvo de 128 MiB, ZSTD nível 3. Rodapé lido antes dos dados.
Grupo de linha64 MiB. Menor unidade descartável por predicado, com mínimo, máximo e nulos por coluna.
models/consumo/fato_pedidos.py
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)
Janela de reprocessamento e custo

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.

dashboards/receita-diaria.json
{
  "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.

src/relatorio.ts
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);
relatorio.py
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ódigoSituaçãoCabeçalho relevanteAção recomendada
200Página retornadaX-A44-Scanned-BytesSeguir com next_cursor até nulo
400Filtro em coluna não publicadaCorrigir; retentar não resolve
403Coluna fora do escopoX-A44-Denied-ColumnsAjustar a projeção ou ampliar o escopo
410Cursor expirado (10 min)Reiniciar a paginação
413Varredura acima da cotaX-A44-Estimated-BytesReduzir datas ou colunas
429Cota esgotadaRetry-AfterBackoff 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.

  1. Identificar o instante alvo

    a44 data history lista versões de manifesto com data, arquivos, linhas e o build que as produziu. Escolha a última anterior ao incidente.

  2. Validar antes de aplicar

    Consulte a versão antiga por viagem no tempo, sem restaurar. Evita restaurar sobre um estado que já estava errado.

  3. 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.

  4. 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.

  5. 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.

terminal
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
Teste de restauração mensal

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.

1 slotaté 180 MB/s WAL
Slots partidos2 a 4 publicações
Origens múltiplasfragmentação por tenant
Limite prático~1,2 GB/s por projeto
DimensãoLimite confortávelTeto observadoO que quebra primeiro depois disso
Taxa de WAL por slot120 MB/s180 MB/sAtraso cresce sem convergir; risco de invalidação
Tabelas por publicação150400Recarga de esquema após DDL passa de 90 s e trava o lote
Volume reescrito por ciclo40 GB150 GBBUILD_OVERLAP; a fila cresce
Partições por tabela20.000120.000Planejamento passa de 5 s só para listar
Modelos no grafo8002.500Resolver dependências domina o build
Consultas por pool1664Derramamento geral; latência 8x pior
Consumidores de API2001.000Avaliar política por requisição vira gargalo de CPU
Escala horizontal não conserta captura

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.

observabilidade/alertas.yaml
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.

tests/consumo.yaml
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
Linhagem de coluna para triagem

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.

sql/politicas.sql
-- 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';
Agregação como canal de vazamento

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ê.

.ci/validar-modelos.sh
#!/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.