PT EN
Voltar ao site

Guia: DATTA Notebook Multi-Fonte

Preview — funcionalidade em desenvolvimento. Comportamento, telas e contratos podem mudar sem aviso entre versões.

A pergunta interessante quase nunca mora numa fonte só: o indicador está no SQL, o relacionamento está no grafo e o contexto está nos documentos. O DATTA Notebook é o ambiente analítico iterativo onde as três convivem — você combina Trino, Neo4j, OpenSearch e PySpark no mesmo arquivo, com as conexões já disponíveis no ambiente, e chega a um resultado reprodutível sem trocar de ferramenta a cada etapa.

Este guia mostra qual fonte responde a qual tipo de pergunta, os padrões de consulta que valem a pena copiar e um exemplo completo de análise de risco integrada.

Qual fonte usar para cada pergunta

FonteLinguagemArmazenamentoTempo típicoMelhor para
Trino 480SQLParquet/ORC em S3/GCS~1–10sAgregações massivas, análises ad-hoc
Neo4jCypherGrafo em memória/disco~100–500msRelacionamentos, linhagem, impacto
OpenSearchLuceneÍndices invertidos~100–500msTexto integral, busca, facetas
PySparkPythonMemória/Parquet~100–500msML, transformações, correlações

1. Como abrir e importar um notebook

O ambiente de notebooks do DATTA é servido pelo JupyterHub.

Opção A: pela interface (recomendado)

  1. Abra http://<endereco-da-plataforma>/jupyter.
  2. Clique em Upload e selecione o arquivo docs/notebook-exemplo-multi-fonte.dattanb.
  3. O notebook é carregado no seu servidor pessoal do JupyterHub e já pode ser executado célula a célula.

Opção B: pela pasta de notebooks

Se o volume de notebooks estiver montado na sua estação de trabalho, basta copiar o arquivo para dentro dele:

bash
cp docs/notebook-exemplo-multi-fonte.dattanb /mnt/notebooks/

Na dúvida sobre o caminho montado na sua instalação, fale com o administrador da plataforma.


2. O que mudou: de Impala para Trino

O motor SQL do notebook é o Trino 480, que substituiu o Impala. Se você tem notebooks antigos, ajuste as consultas ao novo formato de nome de tabela (catalogo.schema.tabela):

Antes (Impala)

sql
-- Impala
SELECT * FROM `default`.`tabela` 
-- Motor: Teradata Impala 4.0.0
-- Armazenamento: GCS gs://datta-object (só Parquet)

Agora (Trino)

sql
-- Trino com catálogo Hive
SELECT * FROM hive.database.tabela
-- Motor: Trino 480
-- Armazenamento: HDFS + S3/GCS (Parquet, ORC, Iceberg)

Catálogos disponíveis

sql
-- Listar catálogos
SHOW CATALOGS;

-- Catálogos do Trino:
-- 1. hive        → Tabelas Hive legacy em Parquet
-- 2. iceberg     → Apache Iceberg (ACID, time-travel)
-- 3. memory      → Tabelas em memória (temp/staging)

3. O que existe em cada fonte

3.1 Trino (Assistência Social)

sql
-- Schema: hive.assistencia_social
-- Tabelas em Parquet com ~23M registros totais

SELECT * FROM hive.assistencia_social.cidadaos;
-- Colunas: id, nome, cpf, nis, data_nascimento, sexo, renda, municipio, bairro, status
-- Registros: ~2.3M

SELECT * FROM hive.assistencia_social.beneficios;
-- Colunas: id, programa, valor, situacao, titular_cpf, titular_nome, banco, data_inicio, tipo
-- Registros: ~5.1M

SELECT * FROM hive.assistencia_social.transacoes;
-- Colunas: id, tipo, valor, data_transacao, programa, conta, municipio, status, cidadao_cpf
-- Registros: ~12.4M

SELECT * FROM hive.assistencia_social.eventos;
-- Colunas: id, tipo, data_evento, descricao, cidadao_cpf, resultado
-- Registros: ~3.8M

SELECT * FROM hive.assistencia_social.alertas;
-- Colunas: id, severidade, tipo, descricao, data_deteccao, regra, status, confianca, cidadao_cpf
-- Registros: ~450K

3.2 Neo4j (grafo — 13 bases)

cypher
-- Base: datta-datacatalog (padrão)
MATCH (n) RETURN labels(n) DISTINCT;
-- Nodes: Dataset, Column, Process, GlossaryTerm, LineageNode

-- Base: datta-ontology
MATCH (n) RETURN labels(n) DISTINCT;
-- Nodes: Ontology, EntityType, PropertyDef, DigitalTwin

-- Base: datta-graph
MATCH (n) RETURN labels(n) DISTINCT;
-- Nodes: Processo, Parte, Empresa, Decisao, Legislacao

3.3 OpenSearch (texto integral)

Três índices de texto estão disponíveis: processos, legislacao e documentos. Você os consulta pelo cliente Python já instalado no ambiente, como mostra o Padrão 3 abaixo.


4. Padrões de uso

Padrão 1: análise distributiva (Trino puro)

Quando usar: agregações massivas, JOINs entre tabelas grandes.

sql
SELECT 
    programa,
    COUNT(*) as qtd_beneficiarios,
    SUM(valor) as valor_total
FROM hive.assistencia_social.beneficios
GROUP BY programa
ORDER BY valor_total DESC;

Desempenho típico: ~2–5 segundos para 5,1M de registros.


Padrão 2: análise semântica (Trino + Neo4j)

Quando usar: relacionamentos, linhagem, impacto.

python
# 1. Extrair dados do Trino (SQL)
df_beneficiarios = spark.sql("""
    SELECT DISTINCT titular_cpf, programa 
    FROM hive.assistencia_social.beneficios
    LIMIT 1000
""")

# 2. Enriquecer com Neo4j (Cypher)
import os
from neo4j import GraphDatabase
driver = GraphDatabase.driver(
    os.environ.get("NEO4J_BOLT_URL", "bolt://neo4j:7687"),
    auth=(os.environ["NEO4J_USERNAME"], os.environ["NEO4J_PASSWORD"])
)
session = driver.session(database="datta-datacatalog")

cpfs = df_beneficiarios.select('titular_cpf').rdd.flatMap(list).collect()
query = """
MATCH (c:Cidadao {cpf: $cpf})-[:ENVOLVIDO_EM]->(p:Processo)
RETURN c.cpf, COUNT(p) as processos
"""
resultados = []
for cpf in cpfs[:100]:  # Sample
    result = session.run(query, cpf=cpf)
    resultados.extend([dict(r) for r in result])

Desempenho típico: ~500ms para 100 consultas individuais — a mesma leitura feita em lote pode ser até 50 vezes mais rápida.


Padrão 3: descoberta contextual (OpenSearch + Trino)

Quando usar: encontrar a legislação relevante para um programa.

python
from elasticsearch import Elasticsearch

es = Elasticsearch(["http://opensearch:9200"])

# Buscar legislação sobre um programa
query = {
    "query": {
        "multi_match": {
            "query": "auxílio emergencial",
            "fields": ["titulo^2", "ementa^1.5", "conteudo"]
        }
    },
    "aggs": {
        "por_ano": {
            "date_histogram": {
                "field": "data_publicacao",
                "calendar_interval": "year"
            }
        }
    }
}

response = es.search(index="legislacao", body=query)

Desempenho típico: ~100–200ms (índice invertido).


5. Exemplo prático: análise de risco integrada

Objetivo

Identificar beneficiários com padrão de risco:

  • renda alta, mas elegível para programa de baixa renda;
  • envolvido em processo judicial questionável;
  • alertas não resolvidos.

Fluxo completo

python
# 1. Trino: Beneficiários anomalias
df_anomalias = spark.sql("""
    SELECT 
        b.titular_cpf,
        b.programa,
        c.renda,
        b.valor,
        COUNT(a.id) as alertas_nao_resolvidos
    FROM hive.assistencia_social.beneficios b
    LEFT JOIN hive.assistencia_social.cidadaos c ON b.titular_cpf = c.cpf
    LEFT JOIN hive.assistencia_social.alertas a ON b.titular_cpf = a.cidadao_cpf 
        AND a.status != 'RESOLVIDO'
    WHERE c.renda > (SELECT PERCENTILE_CONT(0.75) WITHIN GROUP (ORDER BY renda) FROM hive.assistencia_social.cidadaos)
    GROUP BY b.titular_cpf, b.programa, c.renda, b.valor
    HAVING alertas_nao_resolvidos > 0
""")

# 2. Neo4j: Enriquecer com processos
from neo4j import GraphDatabase
driver = GraphDatabase.driver(
    os.environ.get("NEO4J_BOLT_URL", "bolt://neo4j:7687"),
    auth=(os.environ["NEO4J_USERNAME"], os.environ["NEO4J_PASSWORD"])
)
session = driver.session(database="datta-datacatalog")

anomalias_com_processos = []
for row in df_anomalias.collect():
    cpf = row.titular_cpf
    resultado = session.run("""
        MATCH (c:Cidadao {cpf: $cpf})-[:ENVOLVIDO_EM]->(p:Processo)
        WHERE p.status IN ['ATIVO', 'CONTESTADO']
        RETURN COUNT(p) as processos_ativos
    """, cpf=cpf)
    resultado_dict = {**row.asDict(), **dict(resultado.single())}
    anomalias_com_processos.append(resultado_dict)

# 3. OpenSearch: Contexto legislativo
legislacao_relevante = es.search(
    index="legislacao",
    body={
        "query": {
            "terms": {
                "programa": df_anomalias.select('programa').distinct().rdd.flatMap(list).collect()
            }
        }
    }
)

# 4. Pandas: Análise final
import pandas as pd
df_final = pd.DataFrame(anomalias_com_processos)
print(f"Beneficiários em risco: {len(df_final)}")
print(f"Valor total em risco: R$ {df_final['valor'].sum():.2f}")
print(f"Renda média: R$ {df_final['renda'].mean():.2f}")

Ao final você tem um recorte auditável — quem, quanto e por quê — construído a partir de três fontes num único fluxo reprodutível.


6. Desempenho e boas práticas

O que fazer

PadrãoRazão
LIMIT 1000 em exploraçãoEvita output gigante
SELECT col1, col2 (colunas específicas)Parquet é colunar, seleção é rápida
CAST(string_col AS DATE) antes de GROUP BYEvita string parsing N vezes
Reaproveitar valores frequentes no cache da plataformaReduz queries repetidas
spark.sql(...).repartition(200) antes de joinDistribui carga

O que evitar

PadrãoRazão
SELECT *Carrega colunas desnecessárias
JOIN com LIKE '%pattern%'Full scan, sem índice
UDFs Python em loopSerialização cara
collect() sem LIMITPode esgotar a memória da sessão
Queries bloqueantes > 5minExcedem o tempo limite da sessão

7. Quando algo não funciona

SintomaO que fazer
Sessão Spark não respondeNo JupyterHub: Kernel → Restart Kernel e rode as células de novo. Se persistir, peça ao administrador para verificar o estado e os limites de recurso do ambiente de notebooks
No such table: hive.assistencia_social.cidadaosConfirme com SHOW TABLES IN hive.assistencia_social;. Se a lista vier vazia, a base de exemplo ainda não foi carregada — peça ao administrador para rodar o seed do Trino (k8s/trino/seed-trino.sql)
Neo4j connection timeoutConfirme, com os.environ, que as variáveis NEO4J_BOLT_URL, NEO4J_USERNAME e NEO4J_PASSWORD chegaram ao ambiente e que o grafo responde na porta 7687. As credenciais vêm dos segredos da plataforma — nunca escritas no notebook

8. Próximos passos

  1. Executar o notebook de exemplo → entender os padrões.
  2. Criar uma consulta customizada → para o seu caso de uso.
  3. Exportar resultados → para o catálogo (Neo4j ProfilingSnapshot).
  4. Publicar o resultado → o Knowledge Catalog aceita o registro de datasets por integração programática; veja a referência de API.
  5. Agendar jobs → Apache Airflow (opcional).

9. Referências

O ambiente roda sobre JupyterHub, acessível em /jupyter; o motor SQL é o Trino 480 (em substituição ao Impala) e o armazenamento analítico é Parquet/ORC em HDFS + S3/GCS.