Skip to content

Commit 074e28a

Browse files
committed
feat: adiciona ingestão de DAGs de ingestão para empenhos, faturas, cronogramas e terceirizados
1 parent c79aa2e commit 074e28a

4 files changed

Lines changed: 240 additions & 0 deletions

File tree

Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,60 @@
1+
import logging
2+
from airflow.decorators import dag, task
3+
from datetime import datetime, timedelta
4+
from schedule_loader import get_dynamic_schedule
5+
from cliente_contratos import ClienteContratos
6+
from cliente_postgres import ClientPostgresDB
7+
from postgres_helpers import get_postgres_conn
8+
9+
10+
@dag(
11+
dag_id="contratos_cronograma_mir_ingest_dag",
12+
schedule_interval=get_dynamic_schedule("contratos_cronograma_mir_ingest_dag"),
13+
start_date=datetime(2023, 1, 1),
14+
catchup=False,
15+
default_args={
16+
"owner": "Luana",
17+
"retries": 1,
18+
"retry_delay": timedelta(minutes=5),
19+
},
20+
tags=["cronograma_api", "compras_gov", "MIR"],
21+
)
22+
def api_cronograma_mir_dag() -> None:
23+
"""DAG para buscar e armazenar cronograma no PostgreSQL do MIR."""
24+
25+
@task
26+
def fetch_cronograma() -> None:
27+
logging.info("[contratos_cronograma_mir_ingest_dag.py] Starting fetch_cronograma task")
28+
api = ClienteContratos()
29+
postgres_conn_str = get_postgres_conn("postgres_mir")
30+
db = ClientPostgresDB(postgres_conn_str)
31+
contratos_ids = db.get_contratos_ids()
32+
33+
for contrato_id in contratos_ids:
34+
try:
35+
logging.info(
36+
f"[contratos_cronograma_mir_ingest_dag.py] Fetching cronograma for contrato ID: "
37+
f"{contrato_id}"
38+
)
39+
cronograma = api.get_cronograma_by_contrato_id(str(contrato_id))
40+
41+
if cronograma:
42+
for fatura in cronograma:
43+
fatura["dt_ingest"] = datetime.now().isoformat()
44+
45+
db.insert_data(
46+
cronograma,
47+
"cronograma",
48+
conflict_fields=["id"],
49+
primary_key=["id"],
50+
schema="compras_gov",
51+
)
52+
except Exception as e:
53+
logging.error(
54+
f"[contratos_cronograma_mir_ingest_dag.py] Error for contrato ID {contrato_id}: {e}"
55+
)
56+
57+
fetch_cronograma()
58+
59+
60+
dag_instance = api_cronograma_mir_dag()
Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,60 @@
1+
import logging
2+
from airflow.decorators import dag, task
3+
from datetime import datetime, timedelta
4+
from schedule_loader import get_dynamic_schedule
5+
from cliente_contratos import ClienteContratos
6+
from cliente_postgres import ClientPostgresDB
7+
from postgres_helpers import get_postgres_conn
8+
9+
10+
@dag(
11+
dag_id="contratos_empenhos_mir_ingest_dag",
12+
schedule_interval=get_dynamic_schedule("contratos_empenhos_mir_ingest_dag"),
13+
start_date=datetime(2023, 1, 1),
14+
catchup=False,
15+
default_args={
16+
"owner": "Luana",
17+
"retries": 1,
18+
"retry_delay": timedelta(minutes=5),
19+
},
20+
tags=["empenhos_api", "compras_gov", "MIR"],
21+
)
22+
def api_empenhos_mir_dag() -> None:
23+
"""DAG para buscar e armazenar empenhos no PostgreSQL do MIR."""
24+
25+
@task
26+
def fetch_empenhos() -> None:
27+
logging.info("[contratos_empenhos_mir_ingest_dag.py] Starting fetch_empenhos task")
28+
api = ClienteContratos()
29+
postgres_conn_str = get_postgres_conn("postgres_mir")
30+
db = ClientPostgresDB(postgres_conn_str)
31+
contratos_ids = db.get_contratos_ids()
32+
33+
for contrato_id in contratos_ids:
34+
try:
35+
logging.info(
36+
f"[contratos_empenhos_mir_ingest_dag.py] Fetching empenhos for contrato ID: "
37+
f"{contrato_id}"
38+
)
39+
empenhos = api.get_empenhos_by_contrato_id(str(contrato_id))
40+
41+
if empenhos:
42+
for fatura in empenhos:
43+
fatura["dt_ingest"] = datetime.now().isoformat()
44+
45+
db.insert_data(
46+
empenhos,
47+
"empenhos",
48+
conflict_fields=["id"],
49+
primary_key=["id"],
50+
schema="compras_gov",
51+
)
52+
except Exception as e:
53+
logging.error(
54+
f"[contratos_empenhos_mir_ingest_dag.py] Error for contrato ID {contrato_id}: {e}"
55+
)
56+
57+
fetch_empenhos()
58+
59+
60+
dag_instance = api_empenhos_mir_dag()
Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,60 @@
1+
import logging
2+
from airflow.decorators import dag, task
3+
from datetime import datetime, timedelta
4+
from schedule_loader import get_dynamic_schedule
5+
from cliente_contratos import ClienteContratos
6+
from cliente_postgres import ClientPostgresDB
7+
from postgres_helpers import get_postgres_conn
8+
9+
10+
@dag(
11+
dag_id="contratos_faturas_mir_ingest_dag",
12+
schedule_interval=get_dynamic_schedule("contratos_faturas_mir_ingest_dag"),
13+
start_date=datetime(2023, 1, 1),
14+
catchup=False,
15+
default_args={
16+
"owner": "Luana",
17+
"retries": 1,
18+
"retry_delay": timedelta(minutes=5),
19+
},
20+
tags=["faturas_api", "compras_gov", "MIR"],
21+
)
22+
def api_faturas_mir_dag() -> None:
23+
"""DAG para buscar e armazenar faturas no PostgreSQL do MIR."""
24+
25+
@task
26+
def fetch_faturas() -> None:
27+
logging.info("[contratos_faturas_mir_ingest_dag.py] Starting fetch_faturas task")
28+
api = ClienteContratos()
29+
postgres_conn_str = get_postgres_conn("postgres_mir")
30+
db = ClientPostgresDB(postgres_conn_str)
31+
contratos_ids = db.get_contratos_ids()
32+
33+
for contrato_id in contratos_ids:
34+
try:
35+
logging.info(
36+
f"[contratos_faturas_mir_ingest_dag.py] Fetching faturas for contrato ID: "
37+
f"{contrato_id}"
38+
)
39+
faturas = api.get_faturas_by_contrato_id(str(contrato_id))
40+
41+
if faturas:
42+
for fatura in faturas:
43+
fatura["dt_ingest"] = datetime.now().isoformat()
44+
45+
db.insert_data(
46+
faturas,
47+
"faturas",
48+
conflict_fields=["id"],
49+
primary_key=["id"],
50+
schema="compras_gov",
51+
)
52+
except Exception as e:
53+
logging.error(
54+
f"[contratos_faturas_mir_ingest_dag.py] Error for contrato ID {contrato_id}: {e}"
55+
)
56+
57+
fetch_faturas()
58+
59+
60+
dag_instance = api_faturas_mir_dag()
Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,60 @@
1+
import logging
2+
from airflow.decorators import dag, task
3+
from datetime import datetime, timedelta
4+
from schedule_loader import get_dynamic_schedule
5+
from cliente_contratos import ClienteContratos
6+
from cliente_postgres import ClientPostgresDB
7+
from postgres_helpers import get_postgres_conn
8+
9+
10+
@dag(
11+
dag_id="contratos_terceirizados_mir_ingest_dag",
12+
schedule_interval=get_dynamic_schedule("contratos_terceirizados_mir_ingest_dag"),
13+
start_date=datetime(2023, 1, 1),
14+
catchup=False,
15+
default_args={
16+
"owner": "Luana",
17+
"retries": 1,
18+
"retry_delay": timedelta(minutes=5),
19+
},
20+
tags=["terceirizados_api", "compras_gov", "MIR"],
21+
)
22+
def api_terceirizados_mir_dag() -> None:
23+
"""DAG para buscar e armazenar terceirizados no PostgreSQL do MIR."""
24+
25+
@task
26+
def fetch_terceirizados() -> None:
27+
logging.info("[contratos_terceirizados_mir_ingest_dag.py] Starting fetch_terceirizados task")
28+
api = ClienteContratos()
29+
postgres_conn_str = get_postgres_conn("postgres_mir")
30+
db = ClientPostgresDB(postgres_conn_str)
31+
contratos_ids = db.get_contratos_ids()
32+
33+
for contrato_id in contratos_ids:
34+
try:
35+
logging.info(
36+
f"[contratos_terceirizados_mir_ingest_dag.py] Fetching terceirizados for contrato ID: "
37+
f"{contrato_id}"
38+
)
39+
terceirizados = api.get_terceirizados_by_contrato_id(str(contrato_id))
40+
41+
if terceirizados:
42+
for fatura in terceirizados:
43+
fatura["dt_ingest"] = datetime.now().isoformat()
44+
45+
db.insert_data(
46+
terceirizados,
47+
"terceirizados",
48+
conflict_fields=["id"],
49+
primary_key=["id"],
50+
schema="compras_gov",
51+
)
52+
except Exception as e:
53+
logging.error(
54+
f"[contratos_terceirizados_mir_ingest_dag.py] Error for contrato ID {contrato_id}: {e}"
55+
)
56+
57+
fetch_terceirizados()
58+
59+
60+
dag_instance = api_terceirizados_mir_dag()

0 commit comments

Comments
 (0)