PT EN
Voltar ao site

DATTAX — Guia da Linguagem

DATTAX é a linguagem de consulta e transformação do DATTA: uma única sintaxe para ler e cruzar grafo, tabela, índice, arquivo, warehouse e stream. Antes dela, cada resposta exigia dominar um dialeto por motor — SQL no warehouse, Cypher no grafo, consulta própria no índice de busca, mais um framework para eventos ao vivo. Com o DATTAX você escreve um pipeline só, com etapas encadeadas pelo operador |>, e a plataforma traduz cada etapa para a fonte certa — sobra tempo para responder perguntas em vez de trocar de ferramenta.

Este guia apresenta a linguagem do ponto de vista de quem escreve scripts. Ele complementa a referência completa da linguagem, o catálogo de algoritmos de grafo e o guia de streaming.


1. Visão geral

DATTAX é inspirada em pipelines Unix e em linguagens como KQL/PRQL: cada etapa recebe um dataset e produz outro, composto pelo operador |>.

dattax
EVALUATE FROM JDBC "warehouse" "SELECT id, total, created_at FROM vendas"
  |> FILTER total > 0
  |> MUTATE year = YEAR(created_at)
  |> GROUP BY year AGGREGATE total_year = SUM(total)
  |> SELECT year, total_year
  ;

Instruções terminam em ;. Um script pode ter múltiplas instruções.

1.1 Palavras-chave de topo

  • EVALUATE — executa e devolve o dataset ao chamador.
  • DATASET <nome> = <pipeline> ; — nomeia um dataset reutilizável.
  • LET <ident> = <expr> ; — constante local.
  • DEFINE MEASURE <nome> = <expr> ; — medida reutilizável (sum, avg, etc.).
  • MATERIALIZE <pipeline> AS ICEBERG|NEO4J|OPENSEARCH [OPTIONS ...] ; — persiste o resultado no destino escolhido.

1.2 Tipos

  • Escalares: INT, LONG, DOUBLE, STRING, BOOL, TIMESTAMP, DATE.
  • Coluna de grafo: NODE, REL, PATH.
  • Estruturados: LIST<T>, MAP<K,V>.
  • Vetor: VECTOR<dim> (para embeddings / KNN).

Conversões implícitas seguem regras de SQL. Conversão explícita via CAST(x AS INT).


2. De onde vêm os dados (FROM ...)

SintaxeFonte
FROM GRAPH "conn"Neo4j (Cypher + algoritmos de grafo)
FROM INDEX "conn"OpenSearch/Elasticsearch
FROM JDBC "conn" "<sql>"Qualquer banco JDBC catalogado, com o dialeto do fabricante aplicado
FROM TRINO "conn" "<sql>"Trino (federação SQL)
FROM ICEBERG "conn" "<table>"Tabelas Apache Iceberg (catálogo REST)
FROM FILE "<path>"Arquivos no armazenamento da plataforma (CSV/Parquet/JSON/TSV)
FROM STREAM <KIND> "conn"Kafka/Pulsar/Debezium/Neo4j CDC
FROM CASSANDRA "conn" "<cql>"Cassandra/ScyllaDB (CQL)
FROM MONGO "conn" { <pipeline> }MongoDB (aggregation pipeline)
FROM BIGQUERY "conn" "<sql>"Google BigQuery (leitura acelerada pela Storage API)

Todas as conexões são nomes, não URLs. Você cadastra a fonte uma vez em SistemaConexões e passa a referenciá-la pelo nome em qualquer script: ao executar, a plataforma resolve o nome e carrega a credencial daquele usuário no cofre de credenciais. Endereços e senhas nunca aparecem no código.

Exemplos:

dattax
EVALUATE FROM GRAPH "neo4j-main" "MATCH (n:Processo) RETURN n LIMIT 100" ;

EVALUATE FROM INDEX "opensearch-logs" INDEX "app-*"
  QUERY { "match": { "message": "error" } } LIMIT 1000 ;

EVALUATE FROM MONGO "mongo-prod" COLLECTION "orders" {
  { "$match": { "status": "paid" } },
  { "$group": { "_id": "$region", "total": { "$sum": "$amount" } } }
} ;

EVALUATE FROM BIGQUERY "bq-analytics"
  "SELECT country, COUNT(*) n FROM `ds.sessions` GROUP BY country" ;

3. Como transformar os dados

3.1 Básico tabular

OperadorAção
FILTER <expr>Filtro linha a linha (SQL-like)
SELECT col1, col2 AS aliasProjeção / rename
MUTATE col = <expr>Nova coluna calculada (estilo Power Query)
DROP col1, col2Remoção de colunas
RENAME old TO newRenomeação explícita
GROUP BY ... AGGREGATE ...Agregação
`ORDER BY col [ASC\DESC]`
TAKE <n> / SKIP <n>Paginação
DISTINCT [col1, ...]Deduplicação
UNION <dataset>Concatenação vertical

3.2 Join

dattax
|> JOIN outroDs ON a.id = outroDs.id              -- inner
|> LEFT JOIN outroDs ON ...                       -- left outer
|> JOIN outroDs AS LOOKUP ON col = outroDs.key   -- enrichment (O(1))

AS LOOKUP é obrigatório em pipelines de streaming.

3.3 Window

dattax
|> WINDOW PARTITION BY region ORDER BY ts
     RANK() AS rank,
     LAG(valor, 1) AS valor_ant,
     SUM(valor) OVER (ROWS BETWEEN 6 PRECEDING AND CURRENT ROW) AS media_7d

3.4 Pivot / Bucket

dattax
|> PIVOT categoria AGGREGATE total = SUM(valor)
|> BUCKET valor INTO 10                         -- quantile bucketing
|> BUCKET valor INTO [0, 100, 500, 1000]        -- edges explicitas

3.5 Match (subgrafos Cypher)

Quando a fonte é FROM GRAPH, MATCH aceita um pattern Cypher completo:

dattax
EVALUATE FROM GRAPH "neo4j-main"
  |> MATCH (p:Processo)-[:TEM_PARTE]->(parte:Parte {tipo: "autor"})
  |> WHERE p.valor_causa > 10000
  |> RETURN p.numero AS processo, parte.nome AS autor, p.valor_causa
  ;

3.6 KNN (vetor)

Dois usos:

  1. KNN em índice OpenSearch/Elasticsearch (busca vetorial documental, não algoritmo de grafo): ``dattax EVALUATE FROM INDEX "opensearch-docs" |> KNN FIELD "embedding" K 10 VECTOR [0.1, 0.2, ...] ``
  2. KNN em grafo (algoritmo de similaridade): ``dattax |> KNN ON myGraph K=10 ``

4. Algoritmos de grafo (Neo4j Graph Data Science)

Mais de 30 algoritmos do catálogo open source do Neo4j GDS. Todos rodam em modo stream, somente-leitura — um pipeline nunca altera o grafo original. Leiden e SLLPA pertencem à edição Enterprise e estão bloqueados na plataforma: a tentativa de usá-los interrompe a execução com uma mensagem clara em português.

4.1 Centralidade

  • PAGERANK — autoridade global.
  • ARTICLE_RANK — variante que penaliza hubs.
  • EIGENVECTOR — autovetor principal.
  • BETWEENNESS — intermediação.
  • CLOSENESS — proximidade.
  • HARMONIC — closeness harmônica (robusta em grafos desconexos).
  • DEGREE (IN/OUT/BOTH).
  • CELF — influence maximization.

4.2 Comunidade

  • LOUVAIN (via COMMUNITY ALGO=LOUVAIN).
  • LABEL_PROPAGATION.
  • WCC (weakly connected components).
  • SCC (strongly connected components).
  • TRIANGLE_COUNT.
  • LOCAL_CLUSTERING.
  • KCORE.
  • KMEANS (sobre embedding existente).
  • MODULARITY (modularity optimization).

4.3 Pathfinding

  • SHORTEST PATH (Dijkstra).
  • ASTAR (A* com heurística lat/lon).
  • YENS FROM ... TO ... K N (top-K shortest).
  • ALL_SHORTEST_PATHS.
  • BFS FROM ... MAX_DEPTH=N.
  • DFS FROM ....
  • RANDOM_WALK WALK_LENGTH=..., WALKS_PER_NODE=....

4.4 Similaridade

  • NODE_SIMILARITY SIMILARITY_CUTOFF=....
  • FILTERED_NODE_SIMILARITY.
  • KNN K=....
  • FILTERED_KNN K=....

4.5 Embeddings

  • FASTRP EMBEDDING_DIMENSION=..., ITERATIONS=....
  • HASHGNN ITERATIONS=..., EMBEDDING_DIMENSION=....
  • NODE2VEC WALK_LENGTH=..., WALKS_PER_NODE=..., EMBEDDING_DIMENSION=....
  • GRAPHSAGE MODEL "nome" (inferência com modelo pré-treinado).

ADAMIC_ADAR, COMMON_NEIGHBORS, PREFERENTIAL_ATTACHMENT, RESOURCE_ALLOCATION, SAME_COMMUNITY, TOTAL_NEIGHBORS.

dattax
|> LINK_PREDICTION ADAMIC_ADAR BETWEEN "elem-a" AND "elem-b"

Sintaxe e opções de cada algoritmo: catálogo de algoritmos de grafo.

O resultado de cada algoritmo fica em cache por 1 hora, identificado pela combinação de projeção, algoritmo e parâmetros — repetir a mesma execução responde em segundos.


5. Primitivas de IA

Primitivas integradas ao pipeline: a inteligência configurada na plataforma processa texto livre linha a linha, sem você sair do script.

5.1 LLM EXTRACT

Extrai campos estruturados de texto livre.

dattax
EVALUATE FROM JDBC "warehouse" "SELECT id, body FROM tickets"
  |> LLM EXTRACT FROM body INTO {
       categoria: "categoria principal do ticket",
       prioridade: "baixa|media|alta",
       valor_mencionado: "numero em reais se houver, senao null"
     }

5.2 LLM CLASSIFY

Classificação multiclasse com confidence:

dattax
|> LLM CLASSIFY FROM texto LABELS ["spam", "relevante", "duplicado"] AS rotulo

5.3 LLM EMBED

Gera embedding com o modelo da plataforma (Qwen3-Embedding-0.6B por padrão) ou outro configurado:

dattax
|> LLM EMBED FROM descricao AS emb MODEL "Qwen/Qwen3-Embedding-0.6B"

Há quota por usuário e limite de taxa: ao exceder, o script para com mensagem clara em português. Chamadas repetidas sobre o mesmo texto respondem do cache da plataforma, indexado pelo hash do texto.


6. Streaming

Primitivas aditivas — scripts batch não são afetados.

dattax
-- Pedidos ao vivo enriquecidos com usuarios (batch dim):
DATASET users = FROM JDBC "warehouse" "SELECT id, nome, tier FROM users"
  |> CACHE AS LOOKUP REFRESH EVERY "30m" ;

EVALUATE FROM STREAM KAFKA "orders" 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 ;

Primitivas:

  • FROM STREAM KAFKA|PULSAR|CDC|NEO4J_CDC|OPENSEARCH_POLL "<conn>"
  • PARSE TIMESTAMP <col> — marca o event-time.
  • WATERMARK ON <col> DELAY "<dur>" — tolera eventos fora de ordem.
  • WINDOW TUMBLING "<dur>" / SLIDING "<dur>" SLIDE "<dur>" / SESSION "<gap>".
  • CACHE AS LOOKUP [REFRESH EVERY "<dur>"].
  • JOIN <ds> AS LOOKUP ON <eq>.
  • STREAM OUTPUT — terminal.

Detalhes no guia de streaming.


7. Materialização

dattax
MATERIALIZE
  FROM GRAPH "neo4j-main"
  |> MATCH (p:Processo)-[:TEM_PARTE]->(parte)
  |> GROUP BY p.tribunal AGGREGATE qtd = COUNT(*)
AS ICEBERG TABLE "analytics.processos_por_tribunal"
OPTIONS (WRITE_MODE = "OVERWRITE", REFRESH_POLICY = "CRON \"0 0 2 * * *\"") ;

Destinos:

  • ICEBERG — tabelas Apache Iceberg sobre storage compatível com S3 (MinIO).
  • NEO4J — grafo reverso (cria / atualiza labels derivadas).
  • OPENSEARCH — índice analítico.

A plataforma escolhe o plano de escrita (MERGE, OVERWRITE, APPEND) e registra a linhagem no Knowledge Catalog como (:Dataset)-[:DERIVADO_DE]->(:Dataset) — você sempre sabe de onde cada dataset derivou.


8. Medidas (DEFINE MEASURE)

Medidas são expressões nomeadas reutilizáveis:

dattax
DEFINE MEASURE total_vendas = SUM(valor) ;
DEFINE MEASURE ticket_medio = AVG(valor) ;

EVALUATE FROM JDBC "warehouse" "SELECT * FROM vendas"
  |> GROUP BY regiao
     AGGREGATE
       vendas = [total_vendas],
       ticket = [ticket_medio]

9. Exemplos práticos

Para experimentar qualquer exemplo abaixo:

  1. Abra o editor DATTAX da plataforma.
  2. Confirme que as conexões citadas no script existem em SistemaConexões — troque os nomes pelos das suas conexões.
  3. Cole o script, execute e acompanhe o resultado em forma de tabela na própria tela.

9.1 Dashboard analítico (warehouse + grafo)

dattax
-- Processos por tribunal com centralidade de partes
DATASET centralidades =
  FROM GRAPH "neo4j-main"
  |> PROJECT GRAPH g NODES "..." RELS "..."
  |> PAGERANK ON g ITER=20 ;

EVALUATE FROM JDBC "warehouse" "SELECT id, tribunal FROM processos"
  |> JOIN centralidades AS LOOKUP ON id = nodeId
  |> GROUP BY tribunal
     AGGREGATE
       qtd = COUNT(*),
       score_medio = AVG(score)
  |> ORDER BY score_medio DESC
  ;

9.2 Enriquecimento de texto livre com IA

dattax
EVALUATE FROM JDBC "crm" "SELECT id, feedback FROM tickets WHERE data >= CURRENT_DATE - 30"
  |> LLM EXTRACT FROM feedback INTO {
       sentimento: "positivo|neutro|negativo",
       topico: "tema principal",
       nps_inferido: "0-10"
     }
  |> GROUP BY topico, sentimento
     AGGREGATE qtd = COUNT(*)
  ;

9.3 Detecção de fraude em tempo real

dattax
DATASET blacklist = FROM JDBC "compliance" "SELECT cpf FROM lista_negra"
  |> CACHE AS LOOKUP REFRESH EVERY "5m" ;

EVALUATE FROM STREAM KAFKA "transacoes" TOPIC "transacoes.v1"
  |> PARSE TIMESTAMP ts
  |> WATERMARK ON ts DELAY "3s"
  |> JOIN blacklist AS LOOKUP ON cpf_pagador = cpf
  |> FILTER blacklist.cpf IS NOT NULL
  |> WINDOW TUMBLING "1m"
  |> GROUP BY cpf_pagador AGGREGATE qtd = COUNT(*)
  |> FILTER qtd > 3
  |> STREAM OUTPUT ;

9.4 KNN semântico em OpenSearch

dattax
LET query_emb = LLM EMBED_VALUE "contrato inadimplente valor alto" MODEL "Qwen/Qwen3-Embedding-0.6B" ;

EVALUATE FROM INDEX "opensearch-processos"
  |> KNN FIELD "embedding" K 20 VECTOR query_emb
  |> SELECT id, titulo, score
  ;

9.5 Embeddings em lote + materialização

dattax
MATERIALIZE
  FROM JDBC "warehouse" "SELECT id, descricao FROM produtos"
  |> LLM EMBED FROM descricao AS emb MODEL "Qwen/Qwen3-Embedding-0.6B"
AS OPENSEARCH INDEX "produtos-emb"
OPTIONS (MAPPING = "knn_vector:emb:1024") ;

10. Limites e boas práticas

  • Streams não materializáveis: FROM STREAM ... |> MATERIALIZE é rejeitado. Use STREAM OUTPUT + destino externo.
  • IA em volumes grandes: respeite a quota e processe em lotes. LLM EXTRACT em 1M de linhas sem quota para com mensagem clara em português.
  • Algoritmos de grafo em grafos gigantes: use PROJECT GRAPH com um subconjunto via Cypher MATCH ... WHERE .... Evite projeções globais.
  • Cache automático: toda leitura passa por cache (memória local + cache distribuído da plataforma); a invalidação acontece sozinha ao materializar ou alterar a fonte.
  • Scripts grandes: rode EXPLAIN antes de EVALUATE quando disponível.
  • Timeouts: cada fonte tem timeout explícito (padrão 30s). Ajuste com OPTIONS (TIMEOUT = "60s").

11. Para ir além

  • Referência completa da linguagem — gramática, biblioteca padrão, versionamento.
  • Catálogo de algoritmos de grafo — sintaxe e opções de cada algoritmo.
  • Streaming — janelas, watermark e enriquecimento em tempo real.
  • DATTAX × DAX — para quem vem do Power BI.
  • Referência de API — para acionar a plataforma a partir dos seus próprios sistemas.