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 |>.
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 ...)
| Sintaxe | Fonte |
|---|---|
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 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:
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
| Operador | Ação |
|---|---|
FILTER <expr> | Filtro linha a linha (SQL-like) |
SELECT col1, col2 AS alias | Projeção / rename |
MUTATE col = <expr> | Nova coluna calculada (estilo Power Query) |
DROP col1, col2 | Remoção de colunas |
RENAME old TO new | Renomeaçã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
|> 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
|> 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_7d3.4 Pivot / Bucket
|> PIVOT categoria AGGREGATE total = SUM(valor)
|> BUCKET valor INTO 10 -- quantile bucketing
|> BUCKET valor INTO [0, 100, 500, 1000] -- edges explicitas3.5 Match (subgrafos Cypher)
Quando a fonte é FROM GRAPH, MATCH aceita um pattern Cypher completo:
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:
- 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, ...]`` - 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(viaCOMMUNITY 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).
4.6 Link prediction (topológicas)
ADAMIC_ADAR, COMMON_NEIGHBORS, PREFERENTIAL_ATTACHMENT, RESOURCE_ALLOCATION, SAME_COMMUNITY, TOTAL_NEIGHBORS.
|> 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.
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:
|> LLM CLASSIFY FROM texto LABELS ["spam", "relevante", "duplicado"] AS rotulo5.3 LLM EMBED
Gera embedding com o modelo da plataforma (Qwen3-Embedding-0.6B por padrão) ou outro configurado:
|> 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.
-- 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
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:
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:
- Abra o editor DATTAX da plataforma.
- Confirme que as conexões citadas no script existem em — troque os nomes pelos das suas conexões.
- Cole o script, execute e acompanhe o resultado em forma de tabela na própria tela.
9.1 Dashboard analítico (warehouse + grafo)
-- 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
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
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
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
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. UseSTREAM OUTPUT+ destino externo. - IA em volumes grandes: respeite a quota e processe em lotes.
LLM EXTRACTem 1M de linhas sem quota para com mensagem clara em português. - Algoritmos de grafo em grafos gigantes: use
PROJECT GRAPHcom um subconjunto via CypherMATCH ... 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
EXPLAINantes deEVALUATEquando 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.