Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .jules/bolt.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,7 @@
## 2024-10-24 - Streamlit Database Fetch Caching
**Learning:** In Streamlit dashboards, placing `pd.read_sql()` directly in the main script execution path without caching causes the full dataset to be queried from the database and downloaded over the network on every single widget interaction (re-render). This creates a massive performance bottleneck as the data volume grows.
**Action:** Always wrap expensive data fetching operations (like `pd.read_sql`) in Streamlit with `@st.cache_data(ttl=X)` to ensure the data is fetched only once or periodically, making widget interactions lightning fast.

## 2024-05-18 - Pandas iterrows() bottleneck in data loading
**Learning:** Converting dictionaries to Pandas DataFrames only to iterate over them using `df.iterrows()` in an Airflow pipeline creates enormous overhead. Each `iterrows()` iteration yields a new Pandas Series, causing massive performance and memory penalties.
**Action:** Always bypass Pandas DataFrames for row-by-row iteration (e.g. generating tuples for SQL inserts). Iterate directly over the source list of dictionaries using `dict.get()` instead. This can speed up execution by 10x-100x.
54 changes: 25 additions & 29 deletions e2e_open_data_pipeline/dags/public_data_etl.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
import json

# URL de la API de Socrata (Datos Abiertos Colombia)
# Base de datos de origen: "Incidentes viales" (Usemos datos de MedellΓ­n, por ejemplo "yvqg-xvx2" u otra ciudad disponible).
# Base de datos de origen: "Incidentes viales" (Usemos datos de MedellΓ­n, por ejemplo "yvqg-xvx2" u otra ciudad disponible).
# Para este portafolio usaremos "yvqg-xvx2" que corresponde a Secretaria de Movilidad de MedellΓ­n - Incidentes
API_URL = "https://www.datos.gov.co/resource/yvqg-xvx2.json?$limit=1000&$order=fecha%20DESC"

Expand All @@ -22,26 +22,28 @@
'retry_delay': timedelta(minutes=5),
}


def extract_data(**kwargs):
"""Extrae datos de la API pΓΊblica Socrata."""
print(f"Extrayendo datos de: {API_URL}")
response = requests.get(API_URL)
response.raise_for_status()
data = response.json()

# Guardar en XCom para la siguiente tarea
kwargs['ti'].xcom_push(key='raw_data', value=json.dumps(data))
print(f"NΓΊmero de registros obtenidos: {len(data)}")


def transform_data(**kwargs):
"""Transforma el JSON en un DataFrame de Pandas, limpia y normaliza datos."""
ti = kwargs['ti']
raw_data_str = ti.xcom_pull(key='raw_data', task_ids='extract_task')
data = json.loads(raw_data_str)

df = pd.DataFrame(data)
print("Columnas originales:", df.columns)

# Seleccionar y renombrar campos relevantes de la base de MedellΓ­n (Ejemplo)
# Dependiendo de la estructura del JSON, adaptaremos las columnas a la BBDD
column_mapping = {
Expand All @@ -55,17 +57,18 @@ def transform_data(**kwargs):
'latitud': 'latitud',
'longitud': 'longitud'
}

# Intersecar columnas para evitar KeyError si cambian en la API
cols_to_keep = [c for c in column_mapping.keys() if c in df.columns]
df = df[cols_to_keep]
df.rename(columns=column_mapping, inplace=True)

# Limpieza bΓ‘sica
# Convertir a datetime y luego string para postgres
if 'fecha_accidente' in df.columns:
df['fecha_accidente'] = pd.to_datetime(df['fecha_accidente']).dt.date.astype(str)

df['fecha_accidente'] = pd.to_datetime(
df['fecha_accidente']).dt.date.astype(str)

# Filas con lat/long invΓ‘lidos ponerlas como Nulasy luego rellenar a 0 para el mapa (o descartar)
if 'latitud' in df.columns and 'longitud' in df.columns:
df['latitud'] = pd.to_numeric(df['latitud'], errors='coerce')
Expand All @@ -77,33 +80,27 @@ def transform_data(**kwargs):
ti.xcom_push(key='clean_data', value=json.dumps(clean_data))
print(f"Registros despuΓ©s de limpieza: {len(clean_data)}")


def load_data(**kwargs):
"""Carga los datos limpios a PostgreSQL usando ON CONFLICT para evitar exact duplicates."""
ti = kwargs['ti']
clean_data_str = ti.xcom_pull(key='clean_data', task_ids='transform_task')
data = json.loads(clean_data_str)

if not data:
print("No hay datos para cargar.")
return

df = pd.DataFrame(data)

# La conexiΓ³n a BBDD que configuramos en docker compose
# Opcional: configurar Connection Id en la UI de Airflow, usamos 'dw_postgres'
pg_hook = PostgresHook(postgres_conn_id='dw_postgres')

insert_query = """
INSERT INTO public.accidentes_transito (
fecha_accidente, hora_accidente, gravedad_accidente, class_accidente,
lugar_accidente, comuna, barrio, latitud, longitud
) VALUES %s
ON CONFLICT (fecha_accidente, hora_accidente, latitud, longitud) DO NOTHING;
"""

# Preparar records para execute_values

# Preparar records para execute_values iterando directamente sobre la lista de diccionarios
# ⚑ Bolt Optimization: Avoid pd.DataFrame(data) and df.iterrows()
# Iterating directly over the list of dicts avoids creating a pandas Series for every row,
# reducing memory usage and drastically speeding up iteration.
rows = []
for _, row in df.iterrows():
for row in data:
# Usamos .get() con valores default en caso de que alguna columna falte
rows.append((
row.get('fecha_accidente'),
Expand All @@ -116,9 +113,9 @@ def load_data(**kwargs):
row.get('latitud'),
row.get('longitud')
))

from psycopg2.extras import execute_values

# Execute transaction
conn = pg_hook.get_conn()
cursor = conn.cursor()
Expand All @@ -133,39 +130,38 @@ def load_data(**kwargs):
"""
execute_values(cursor, insert_query, rows)
conn.commit()
print(f"Carga exitosa! Se han insertado/ignorado {len(rows)} registros.")
print(
f"Carga exitosa! Se han insertado/ignorado {len(rows)} registros.")
except Exception as e:
conn.rollback()
raise e
finally:
cursor.close()
conn.close()


with DAG(
'open_data_etl_accidentes',
default_args=default_args,
description='ETL pipeline para obtener accidentes de trΓ‘nsito de Datos Abiertos Colombia',
schedule_interval=timedelta(days=1),
schedule=timedelta(days=1),
catchup=False,
tags=['portafolio', 'datos_abiertos'],
) as dag:

extract_task = PythonOperator(
task_id='extract_task',
python_callable=extract_data,
provide_context=True,
)

transform_task = PythonOperator(
task_id='transform_task',
python_callable=transform_data,
provide_context=True,
)

load_task = PythonOperator(
task_id='load_task',
python_callable=load_data,
provide_context=True,
)

extract_task >> transform_task >> load_task
97 changes: 97 additions & 0 deletions e2e_open_data_pipeline/test_load_data.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
import json
import pytest
from unittest.mock import MagicMock, patch
import pandas as pd
import datetime

# Import the functions to test
from e2e_open_data_pipeline.dags.public_data_etl import extract_data, transform_data, load_data


@pytest.fixture
def mock_kwargs():
ti_mock = MagicMock()
return {'ti': ti_mock}


@patch('requests.get')
def test_extract_data(mock_get, mock_kwargs):
# Mock API response
mock_response = MagicMock()
mock_response.json.return_value = [
{'fecha': '2023-01-01T00:00:00.000', 'latitud': '6.2', 'longitud': '-75.5'}]
mock_get.return_value = mock_response

extract_data(**mock_kwargs)

# Verify xcom_push was called with the correct data
mock_kwargs['ti'].xcom_push.assert_called_once()
args, kwargs = mock_kwargs['ti'].xcom_push.call_args
assert kwargs['key'] == 'raw_data'
assert json.loads(kwargs['value']) == [
{'fecha': '2023-01-01T00:00:00.000', 'latitud': '6.2', 'longitud': '-75.5'}]


def test_transform_data(mock_kwargs):
# Mock xcom_pull
raw_data = [{'fecha': '2023-01-01T00:00:00.000', 'hora': '10:00:00',
'latitud': '6.2', 'longitud': '-75.5', 'gravedad': 'HERIDO'}]
mock_kwargs['ti'].xcom_pull.return_value = json.dumps(raw_data)

transform_data(**mock_kwargs)

# Verify xcom_push
mock_kwargs['ti'].xcom_push.assert_called_once()
args, kwargs = mock_kwargs['ti'].xcom_push.call_args
assert kwargs['key'] == 'clean_data'

clean_data = json.loads(kwargs['value'])
assert len(clean_data) == 1
assert clean_data[0]['fecha_accidente'] == '2023-01-01'
assert clean_data[0]['hora_accidente'] == '10:00:00'
assert clean_data[0]['latitud'] == 6.2
assert clean_data[0]['longitud'] == -75.5
assert clean_data[0]['gravedad_accidente'] == 'HERIDO'


@patch('e2e_open_data_pipeline.dags.public_data_etl.PostgresHook')
@patch('psycopg2.extras.execute_values')
def test_load_data(mock_execute_values, mock_postgres_hook, mock_kwargs):
# Mock xcom_pull
clean_data = [{
'fecha_accidente': '2023-01-01',
'hora_accidente': '10:00:00',
'gravedad_accidente': 'HERIDO',
'clase_accidente': 'CHOQUE',
'lugar_accidente': 'CALLE 1',
'comuna': 'CANDELARIA',
'barrio': 'CENTRO',
'latitud': 6.2,
'longitud': -75.5
}]
mock_kwargs['ti'].xcom_pull.return_value = json.dumps(clean_data)

# Setup mock PostgresHook
mock_conn = MagicMock()
mock_cursor = MagicMock()
mock_conn.cursor.return_value = mock_cursor
mock_postgres_hook.return_value.get_conn.return_value = mock_conn

load_data(**mock_kwargs)

# Verify execution
mock_postgres_hook.assert_called_once_with(postgres_conn_id='dw_postgres')
mock_execute_values.assert_called_once()

args, kwargs = mock_execute_values.call_args
assert args[0] == mock_cursor
# Check the actual values passed in the rows list
rows_passed = args[2]
assert len(rows_passed) == 1
assert rows_passed[0] == (
'2023-01-01', '10:00:00', 'HERIDO', 'CHOQUE', 'CALLE 1', 'CANDELARIA', 'CENTRO', 6.2, -75.5
)

mock_conn.commit.assert_called_once()
mock_cursor.close.assert_called_once()
mock_conn.close.assert_called_once()