DATTAX Streaming — Sintaxe
Acompanhar dados em tempo real costumava exigir um projeto de engenharia à parte: consumidores Kafka, gerência de offsets, janelas, estado. No DATTAX, streaming é a mesma linguagem dos seus scripts batch — você troca o FROM por FROM STREAM, adiciona janela e watermark, e o pipeline passa a emitir resultados ao vivo. Todas as primitivas são aditivas: scripts batch existentes continuam funcionando sem alteração.
Fonte: FROM STREAM ...
FROM STREAM <KIND> "<connectionName>" [options]Tipos suportados:
| Kind | Descrição |
|---|---|
KAFKA | Tópico Kafka. |
PULSAR | Tópico Pulsar. |
CDC | Tópico Kafka populado por Debezium Connect. |
NEO4J_CDC | Procedure cdc.query do Neo4j, com leitura periódica como alternativa. |
OPENSEARCH_POLL | Pseudo-stream sobre OpenSearch (point-in-time + consulta periódica). |
Opções (todas opcionais, na ordem que aparecem)
TOPIC "nome" -- obrigatorio para KAFKA/PULSAR/OPENSEARCH_POLL
TABLE "schema.tabela" -- obrigatorio para CDC
LABEL "Label" -- obrigatorio para NEO4J_CDC fallback
FORMAT JSON | AVRO | PROTOBUF
STARTING FROM EARLIEST | LATEST | TIMESTAMP "2026-01-01T00:00:00Z"
GROUP ID "meu-consumer" -- default: datta-dattax-<subscriptionId>Exemplos:
EVALUATE FROM STREAM KAFKA "my-kafka" TOPIC "orders" FORMAT JSON STARTING FROM LATEST ;
EVALUATE FROM STREAM CDC "postgres-prd" TABLE "public.orders" ;
EVALUATE FROM STREAM NEO4J_CDC "neo4j-main" LABEL "Order" STARTING FROM TIMESTAMP "2026-01-01T00:00:00Z" ;Importante:
connectionNameé o identificador lógico de uma conexão já cadastrada na plataforma — veja o guia de conexões. Endereços de servidores e credenciais nunca aparecem no script: ficam no cofre seguro da plataforma.
Transformações de stream
PARSE TIMESTAMP <column>
Marca a coluna como event-time — requisito para watermark/window.
|> PARSE TIMESTAMP tsWATERMARK ON <column> DELAY "<duration>"
Tolera eventos fora de ordem até delay depois do timestamp.
|> WATERMARK ON ts DELAY "5s"WINDOW ...
Três tipos:
|> WINDOW TUMBLING "30s"
|> WINDOW SLIDING "60s" SLIDE "10s"
|> WINDOW SESSION "5m"A emissão ocorre quando o watermark passa o fim da janela (tumbling/sliding) ou quando não chegam eventos dentro do gap (session).
CACHE AS LOOKUP [REFRESH EVERY "<duration>"]
Materializa um dataset batch em memória para enriquecimento:
DATASET users = FROM JDBC "warehouse" "SELECT id, nome, tier FROM users"
|> CACHE AS LOOKUP REFRESH EVERY "30m" ;JOIN ... AS LOOKUP
Enriquecimento unilateral: cada evento do stream busca sua chave no lookup em tempo constante — sem atrasar o fluxo.
EVALUATE FROM STREAM KAFKA "my-kafka" TOPIC "events"
|> JOIN users AS LOOKUP ON user_id = id
|> STREAM OUTPUT ;JOIN ... AS STREAM
Join entre dois streams ainda não é suportado — a tentativa exibe uma mensagem em português sugerindo AS LOOKUP. Previsto para uma versão futura.
STREAM OUTPUT
Marca o pipeline como terminal streaming. A assinatura publica eventos em tempo real na interface até ser cancelada.
Exemplo prático — pedidos premium ao vivo
Cenário: você quer ver, a cada 30 segundos, os pedidos de clientes premium chegando pelo tópico Kafka orders, já enriquecidos com o cadastro de usuários do warehouse.
- Confirme que as conexões
orders-stream(Kafka) ewarehouse(JDBC) já estão cadastradas na plataforma. - No editor DATTAX, cole e execute:
-- Enrichment de pedidos ao vivo com dim de usuarios
DATASET users = FROM JDBC "warehouse" "SELECT id, nome, tier FROM users"
|> CACHE AS LOOKUP REFRESH EVERY "30m" ;
EVALUATE FROM STREAM KAFKA "orders-stream" TOPIC "orders" STARTING FROM LATEST
|> PARSE TIMESTAMP event_ts
|> WATERMARK ON event_ts DELAY "5s"
|> JOIN users AS LOOKUP ON user_id = id
|> FILTER tier = "premium"
|> WINDOW TUMBLING "30s"
|> STREAM OUTPUT ;- A cada janela fechada, novas linhas aparecem na tela. A assinatura segue ativa — e listada no painel de execuções em segundo plano — até você cancelá-la.
Assinaturas de stream também podem ser criadas, acompanhadas (inclusive o atraso de consumo) e canceladas por integração: a plataforma expõe endpoints para criar a assinatura, receber os eventos por WebSocket ou pelo canal de eventos alternativo, consultar as métricas de atraso e encerrar a assinatura — veja a referência de API.
Permissões
DATTABI_STREAM_VIEW— assinar/ver (analista+)DATTABI_STREAM_CREATE— criar assinatura (usuário avançado+)DATTABI_STREAM_ADMIN— gerenciar conectores de CDC (admin)