PT EN
Voltar ao site

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:

KindDescrição
KAFKATópico Kafka.
PULSARTópico Pulsar.
CDCTópico Kafka populado por Debezium Connect.
NEO4J_CDCProcedure cdc.query do Neo4j, com leitura periódica como alternativa.
OPENSEARCH_POLLPseudo-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 ts

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

  1. Confirme que as conexões orders-stream (Kafka) e warehouse (JDBC) já estão cadastradas na plataforma.
  2. 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 ;
  1. 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)