Skip to content

Commit f9e407a

Browse files
authored
Merge pull request #197 from GovHub-br/feat/ingestao-siconv
feat: adiciona dag de ingestão de dados do SICONV
2 parents 12b432f + 854b650 commit f9e407a

3 files changed

Lines changed: 375 additions & 0 deletions

File tree

Lines changed: 111 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,111 @@
1+
import logging
2+
from datetime import datetime, timedelta
3+
from airflow.decorators import dag, task
4+
from postgres_helpers import get_postgres_conn
5+
from cliente_postgres import ClientPostgresDB
6+
from cliente_siconv import ClienteSiconv
7+
from tabelas_siconv import TABELAS_SICONV
8+
9+
@dag(
10+
schedule_interval= None,
11+
start_date=datetime(2024, 1, 1),
12+
catchup=False,
13+
default_args={
14+
"owner": "Luana",
15+
"retries": 1,
16+
"retry_delay": timedelta(minutes=5),
17+
},
18+
tags=["siconv", "MIR"],
19+
)
20+
def siconv_ingestao_dag() -> None:
21+
22+
@task
23+
def baixar_siconv() -> str:
24+
cliente = ClienteSiconv()
25+
cliente.baixar_zip()
26+
return cliente.ZIP_PATH
27+
28+
@task
29+
def ingerir_tabela(zip_path: str, nome_tabela: str, nome_csv: str, conflict_fields: list, primary_key: list, skip_rows: int, colunas: list, truncate_before_insert: bool = False) -> None:
30+
logging.info(f"[siconv_ingest_dag.py] Iniciando ingestão da tabela {nome_tabela}")
31+
postgres_conn_str = get_postgres_conn("postgres_mir")
32+
db = ClientPostgresDB(postgres_conn_str)
33+
cliente = ClienteSiconv()
34+
35+
if truncate_before_insert:
36+
logging.info(f"[siconv_ingest_dag.py] Truncando tabela siconv.{nome_tabela}...")
37+
db.execute_non_query(f"""
38+
DO $$ BEGIN
39+
IF EXISTS (SELECT FROM pg_tables WHERE schemaname = 'siconv' AND tablename = '{nome_tabela}') THEN
40+
TRUNCATE TABLE siconv.{nome_tabela};
41+
END IF;
42+
END $$;
43+
""")
44+
gerador_registros = cliente.ler_csv(nome_csv, skip_rows, colunas_esperadas=colunas)
45+
46+
lote = []
47+
tamanho_lote = 5000
48+
total_inserido = 0
49+
50+
for registro in gerador_registros:
51+
lote.append(registro)
52+
53+
if len(lote) >= tamanho_lote:
54+
lote = [dict(t) for t in {tuple(d.items()) for d in lote}]
55+
db.insert_data(
56+
lote,
57+
nome_tabela,
58+
conflict_fields=conflict_fields,
59+
primary_key=primary_key,
60+
schema="siconv",
61+
)
62+
total_inserido += len(lote)
63+
logging.info(f"[siconv_ingest_dag.py] {total_inserido} registros processados...")
64+
lote = []
65+
66+
if lote:
67+
lote = [dict(t) for t in {tuple(d.items()) for d in lote}]
68+
db.insert_data(
69+
lote,
70+
nome_tabela,
71+
conflict_fields=conflict_fields,
72+
primary_key=primary_key,
73+
schema="siconv",
74+
)
75+
total_inserido += len(lote)
76+
77+
if total_inserido == 0:
78+
logging.warning(f"[siconv_ingest_dag.py] Nenhum registro processado para {nome_tabela}")
79+
else:
80+
logging.info(f"[siconv_ingest_dag.py] Ingestão finalizada: {total_inserido} registros em {nome_tabela}")
81+
82+
@task
83+
def deletar_zip(zip_path: str) -> None:
84+
import os
85+
if os.path.exists(zip_path):
86+
os.remove(zip_path)
87+
logging.info(f"[siconv_ingest_dag.py] Arquivo {zip_path} deletado com sucesso")
88+
else:
89+
logging.warning(f"[siconv_ingest_dag.py] Arquivo {zip_path} não encontrado para deletar")
90+
91+
path_zip = baixar_siconv()
92+
93+
ultima_task = path_zip
94+
95+
for tabela in TABELAS_SICONV:
96+
task_atual = ingerir_tabela.override(task_id=f"ingerir_{tabela['nome_tabela']}")(
97+
zip_path=path_zip,
98+
nome_tabela=tabela["nome_tabela"],
99+
nome_csv=tabela["nome_csv"],
100+
conflict_fields=tabela["conflict_fields"],
101+
primary_key=tabela["primary_key"],
102+
skip_rows=tabela["skip_rows"],
103+
colunas=tabela["colunas"],
104+
truncate_before_insert=tabela.get("truncate_before_insert", False),
105+
)
106+
107+
ultima_task >> task_atual
108+
ultima_task = task_atual
109+
ultima_task >> deletar_zip(path_zip)
110+
111+
siconv_ingestao_dag()
Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,40 @@
1+
import zipfile
2+
import csv
3+
import io
4+
import logging
5+
import requests
6+
7+
class ClienteSiconv:
8+
URL_ZIP = "https://repositorio.dados.gov.br/seges/detru/siconv.zip"
9+
ZIP_PATH = "/tmp/siconv.zip"
10+
11+
def baixar_zip(self) -> None:
12+
logging.info("[cliente_siconv.py] Baixando arquivo SICONV...")
13+
response = requests.get(self.URL_ZIP, stream=True)
14+
response.raise_for_status()
15+
with open(self.ZIP_PATH, "wb") as f:
16+
for chunk in response.iter_content(chunk_size=8192):
17+
f.write(chunk)
18+
logging.info("[cliente_siconv.py] Download concluído")
19+
20+
def ler_csv(self, nome_csv: str, skip_rows: int = 0, colunas_esperadas: list = None):
21+
logging.info(f"[cliente_siconv.py] Lendo {nome_csv} em modo streaming...")
22+
with zipfile.ZipFile(self.ZIP_PATH, "r") as z:
23+
with z.open(nome_csv) as f:
24+
conteudo = io.TextIOWrapper(f, encoding="utf-8-sig")
25+
reader = csv.DictReader(conteudo, delimiter=";")
26+
27+
if colunas_esperadas:
28+
colunas_csv = reader.fieldnames or []
29+
faltando = [c for c in colunas_esperadas if c not in colunas_csv]
30+
if faltando:
31+
raise ValueError(f"[cliente_siconv.py] Colunas faltando em {nome_csv}: {faltando}")
32+
33+
for i, row in enumerate(reader):
34+
if i < skip_rows:
35+
continue
36+
37+
if colunas_esperadas:
38+
yield {k.lower(): row[k] for k in colunas_esperadas}
39+
else:
40+
yield {k.lower(): v for k, v in row.items() if k is not None}
Lines changed: 224 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,224 @@
1+
TABELAS_SICONV = [
2+
{
3+
"nome_tabela": "proposta",
4+
"nome_csv": "siconv_proposta.csv",
5+
"conflict_fields": ["id_proposta"],
6+
"primary_key": ["id_proposta"],
7+
"truncate_before_insert": False,
8+
"skip_rows": 0,
9+
"colunas": [
10+
"ID_PROPOSTA", "UF_PROPONENTE", "MUNIC_PROPONENTE", "COD_MUNIC_IBGE",
11+
"COD_ORGAO_SUP", "DESC_ORGAO_SUP", "NATUREZA_JURIDICA", "NR_PROPOSTA",
12+
"DIA_PROP", "MES_PROP", "ANO_PROP", "DIA_PROPOSTA", "COD_ORGAO",
13+
"DESC_ORGAO", "MODALIDADE", "IDENTIF_PROPONENTE", "NM_PROPONENTE",
14+
"CEP_PROPONENTE", "ENDERECO_PROPONENTE", "BAIRRO_PROPONENTE", "NM_BANCO",
15+
"SITUACAO_CONTA", "SITUACAO_PROJETO_BASICO", "SIT_PROPOSTA",
16+
"DIA_INIC_VIGENCIA_PROPOSTA", "DIA_FIM_VIGENCIA_PROPOSTA",
17+
"OBJETO_PROPOSTA", "ITEM_INVESTIMENTO", "ENVIADA_MANDATARIA",
18+
"NOME_SUBTIPO_PROPOSTA", "DESCRICAO_SUBTIPO_PROPOSTA", "VL_GLOBAL_PROP",
19+
"VL_REPASSE_PROP", "VL_CONTRAPARTIDA_PROP", "CD_AGENCIA", "CD_CONTA",
20+
],
21+
},
22+
{
23+
"nome_tabela": "convenio",
24+
"nome_csv": "siconv_convenio.csv",
25+
"conflict_fields": ["nr_convenio"],
26+
"primary_key": ["nr_convenio"],
27+
"truncate_before_insert": False,
28+
"skip_rows": 0,
29+
"colunas": [
30+
"NR_CONVENIO", "ID_PROPOSTA", "DIA", "MES", "ANO", "DIA_ASSIN_CONV",
31+
"SIT_CONVENIO", "SUBSITUACAO_CONV", "SITUACAO_PUBLICACAO", "INSTRUMENTO_ATIVO",
32+
"IND_OPERA_OBTV", "NR_PROCESSO", "UG_EMITENTE", "DIA_PUBL_CONV",
33+
"DIA_INIC_VIGENC_CONV", "DIA_FIM_VIGENC_CONV", "DIA_FIM_VIGENC_ORIGINAL_CONV",
34+
"DIAS_PREST_CONTAS", "DIA_LIMITE_PREST_CONTAS", "DATA_SUSPENSIVA",
35+
"DATA_RETIRADA_SUSPENSIVA", "DIAS_CLAUSULA_SUSPENSIVA", "SITUACAO_CONTRATACAO",
36+
"IND_ASSINADO", "MOTIVO_SUSPENSAO", "IND_FOTO", "QTDE_CONVENIOS", "QTD_TA",
37+
"QTD_PRORROGA", "VL_GLOBAL_CONV", "VL_REPASSE_CONV", "VL_CONTRAPARTIDA_CONV",
38+
"VL_EMPENHADO_CONV", "VL_DESEMBOLSADO_CONV", "VL_SALDO_REMAN_TESOURO",
39+
"VL_SALDO_REMAN_CONVENENTE", "VL_RENDIMENTO_APLICACAO", "VL_INGRESSO_CONTRAPARTIDA",
40+
"VL_SALDO_CONTA", "VALOR_GLOBAL_ORIGINAL_CONV",
41+
],
42+
},
43+
{
44+
"nome_tabela": "desembolso",
45+
"nome_csv": "siconv_desembolso.csv",
46+
"conflict_fields": ["id_desembolso"],
47+
"primary_key": ["id_desembolso"],
48+
"truncate_before_insert": False,
49+
"skip_rows": 0,
50+
"colunas": [
51+
"ID_DESEMBOLSO", "NR_CONVENIO", "DT_ULT_DESEMBOLSO", "QTD_DIAS_SEM_DESEMBOLSO",
52+
"DATA_DESEMBOLSO", "ANO_DESEMBOLSO", "MES_DESEMBOLSO", "NR_SIAFI",
53+
"UG_EMITENTE_DH", "OBSERVACAO_DH", "VL_DESEMBOLSADO",
54+
],
55+
},
56+
{
57+
"nome_tabela": "desbloqueio",
58+
"nome_csv": "siconv_desbloqueio_cr.csv",
59+
"conflict_fields": [],
60+
"primary_key": [],
61+
"truncate_before_insert": True,
62+
"skip_rows": 0,
63+
"colunas": [
64+
"NR_CONVENIO", "NR_OB", "DATA_CADASTRO", "DATA_ENVIO",
65+
"TIPO_RECURSO_DESBLOQUEIO", "VL_TOTAL_DESBLOQUEIO",
66+
"VL_DESBLOQUEADO", "VL_BLOQUEADO",
67+
],
68+
},
69+
{
70+
"nome_tabela": "solicitacao_alteracao",
71+
"nome_csv": "siconv_solicitacao_alteracao.csv",
72+
"conflict_fields": ["id_solicitacao"],
73+
"primary_key": ["id_solicitacao"],
74+
"truncate_before_insert": False,
75+
"skip_rows": 0,
76+
"colunas": [
77+
"ID_SOLICITACAO", "NR_CONVENIO", "NR_SOLICITACAO", "SITUACAO_SOLICITACAO",
78+
"OBJETO_SOLICITACAO", "DATA_SOLICITACAO",
79+
],
80+
},
81+
{
82+
"nome_tabela": "termo_aditivo",
83+
"nome_csv": "siconv_termo_aditivo.csv",
84+
"conflict_fields": ["nr_convenio", "numero_ta"],
85+
"primary_key": ["nr_convenio", "numero_ta"],
86+
"truncate_before_insert": False,
87+
"skip_rows": 0,
88+
"colunas": [
89+
"NR_CONVENIO", "ID_SOLICITACAO", "NUMERO_TA", "TIPO_TA",
90+
"VL_GLOBAL_TA", "VL_REPASSE_TA", "VL_CONTRAPARTIDA_TA",
91+
"DT_ASSINATURA_TA", "DT_INICIO_TA", "DT_FIM_TA", "JUSTIFICATIVA_TA",
92+
],
93+
},
94+
{
95+
"nome_tabela": "solicitacao_rendimento_aplicacao",
96+
"nome_csv": "siconv_solicitacao_rendimento_aplicacao.csv",
97+
"conflict_fields": ["id_solicitacao_rend_aplicacao"],
98+
"primary_key": ["id_solicitacao_rend_aplicacao"],
99+
"truncate_before_insert": False,
100+
"skip_rows": 0,
101+
"colunas": [
102+
"ID_SOLICITACAO_REND_APLICACAO", "NR_CONVENIO", "NR_SOLICITACAO_REND_APLICACAO",
103+
"STATUS_SOLICITACAO_REND_APLICACAO", "DATA_SOLICITACAO_REND_APLICACAO",
104+
"VALOR_SOLICITACAO_REND_APLICACAO", "VALOR_APROVADO_SOLICITACAO_REND_APLICACAO",
105+
],
106+
},
107+
{
108+
"nome_tabela": "prorroga_oficio",
109+
"nome_csv": "siconv_prorroga_oficio.csv",
110+
"conflict_fields": [],
111+
"primary_key": [],
112+
"truncate_before_insert": True,
113+
"skip_rows": 0,
114+
"colunas": [
115+
"NR_CONVENIO", "NR_PRORROGA", "DT_INICIO_PRORROGA", "DT_FIM_PRORROGA",
116+
"DIAS_PRORROGA", "DT_ASSINATURA_PRORROGA", "SIT_PRORROGA",
117+
],
118+
},
119+
{
120+
"nome_tabela": "pagamento_tributo",
121+
"nome_csv": "siconv_pagamento_tributo.csv",
122+
"conflict_fields": ["nr_convenio", "data_tributo"],
123+
"primary_key": ["nr_convenio", "data_tributo"],
124+
"truncate_before_insert": False,
125+
"skip_rows": 0,
126+
"colunas": [
127+
"NR_CONVENIO", "DATA_TRIBUTO", "VL_PAG_TRIBUTOS",
128+
],
129+
},
130+
{
131+
"nome_tabela": "pagamento",
132+
"nome_csv": "siconv_pagamento.csv",
133+
"conflict_fields": ["nr_mov_fin"],
134+
"primary_key": ["nr_mov_fin"],
135+
"truncate_before_insert": False,
136+
"skip_rows": 0,
137+
"colunas": [
138+
"NR_MOV_FIN", "NR_CONVENIO", "IDENTIF_FORNECEDOR", "NOME_FORNECEDOR",
139+
"TP_MOV_FINANCEIRA", "DATA_PAG", "NR_DL", "DESC_DL",
140+
"VL_PAGO", "ID_DL", "DATA_EMISSAO_DL",
141+
],
142+
},
143+
{
144+
"nome_tabela": "licitacao",
145+
"nome_csv": "siconv_licitacao.csv",
146+
"conflict_fields": ["id_licitacao"],
147+
"primary_key": ["id_licitacao"],
148+
"truncate_before_insert": False,
149+
"skip_rows": 0,
150+
"colunas": [
151+
"ID_LICITACAO", "NR_CONVENIO", "NR_LICITACAO", "MODALIDADE_LICITACAO",
152+
"TP_PROCESSO_COMPRA", "TIPO_LICITACAO", "NR_PROCESSO_LICITACAO",
153+
"DATA_PUBLICACAO_LICITACAO", "DATA_ABERTURA_LICITACAO", "DATA_ENCERRAMENTO_LICITACAO",
154+
"DATA_HOMOLOGACAO_LICITACAO", "STATUS_LICITACAO", "SITUACAO_ACEITE_PROCESSO_EXECU",
155+
"SISTEMA_ORIGEM", "SITUACAO_SISTEMA", "VALOR_LICITACAO",
156+
"DATA_ANALISE_ACEITE", "DATA_ENVIO_ANALISE",
157+
],
158+
},
159+
{
160+
"nome_tabela": "ingresso_contrapartida",
161+
"nome_csv": "siconv_ingresso_contrapartida.csv",
162+
"conflict_fields": ["nr_convenio", "dt_ingresso_contrapartida"],
163+
"primary_key": ["nr_convenio", "dt_ingresso_contrapartida"],
164+
"truncate_before_insert": False,
165+
"skip_rows": 0,
166+
"colunas": [
167+
"NR_CONVENIO", "DT_INGRESSO_CONTRAPARTIDA", "VL_INGRESSO_CONTRAPARTIDA",
168+
],
169+
},
170+
{
171+
"nome_tabela": "empenho",
172+
"nome_csv": "siconv_empenho.csv",
173+
"conflict_fields": [],
174+
"primary_key": [],
175+
"truncate_before_insert": True,
176+
"skip_rows": 0,
177+
"colunas": [
178+
"ID_EMPENHO", "NR_CONVENIO", "NR_EMPENHO", "TIPO_NOTA", "DESC_TIPO_NOTA",
179+
"DATA_EMISSAO", "COD_SITUACAO_EMPENHO", "DESC_SITUACAO_EMPENHO",
180+
"UG_EMITENTE", "UG_RESPONSAVEL", "FONTE_RECURSO", "NATUREZA_DESPESA",
181+
"PLANO_INTERNO", "PTRES", "VALOR_EMPENHO", "RESULTADO_PRIMARIO",
182+
"OBSERVACAO_EMPENHO", "DESCRICAO_EMENDA_SIAFI",
183+
],
184+
},
185+
{
186+
"nome_tabela": "historico_situacao",
187+
"nome_csv": "siconv_historico_situacao.csv",
188+
"conflict_fields": [],
189+
"primary_key": [],
190+
"truncate_before_insert": True,
191+
"skip_rows": 0,
192+
"colunas": [
193+
"ID_PROPOSTA", "NR_CONVENIO", "DIA_HISTORICO_SIT", "HISTORICO_SIT",
194+
"DIAS_HISTORICO_SIT", "COD_HISTORICO_SIT",
195+
],
196+
},
197+
{
198+
"nome_tabela": "cronograma_desembolso",
199+
"nome_csv": "siconv_cronograma_desembolso.csv",
200+
"conflict_fields": ["id_proposta", "nr_convenio", "nr_parcela_crono_desembolso"],
201+
"primary_key": ["id_proposta", "nr_convenio", "nr_parcela_crono_desembolso"],
202+
"truncate_before_insert": False,
203+
"skip_rows": 0,
204+
"colunas": [
205+
"ID_PROPOSTA", "NR_CONVENIO", "NR_PARCELA_CRONO_DESEMBOLSO",
206+
"MES_CRONO_DESEMBOLSO", "ANO_CRONO_DESEMBOLSO", "TIPO_RESP_CRONO_DESEMBOLSO",
207+
"VALOR_PARCELA_CRONO_DESEMBOLSO",
208+
],
209+
},
210+
{
211+
"nome_tabela": "meta_crono_fisico",
212+
"nome_csv": "siconv_meta_crono_fisico.csv",
213+
"conflict_fields": ["id_meta"],
214+
"primary_key": ["id_meta"],
215+
"truncate_before_insert": False,
216+
"skip_rows": 0,
217+
"colunas": [
218+
"ID_META", "ID_PROPOSTA", "NR_CONVENIO", "COD_PROGRAMA", "NOME_PROGRAMA",
219+
"NR_META", "TIPO_META", "DESC_META", "DATA_INICIO_META", "DATA_FIM_META",
220+
"UF_META", "MUNICIPIO_META", "ENDERECO_META", "CEP_META",
221+
"QTD_META", "UND_FORNECIMENTO_META", "VL_META",
222+
],
223+
},
224+
]

0 commit comments

Comments
 (0)