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
3 changes: 3 additions & 0 deletions .jules/bolt.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,6 @@
## 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.
## 2025-02-12 - Pandas iterrows() Optimization
**Learning:** Pandas `iterrows()` is notoriously slow because it converts each row to a Series object, creating significant overhead, especially during database loads. This is a common performance bottleneck in ETL processes handling moderate to large datasets.
**Action:** When extracting data from a list of dictionaries (like JSON API responses) to insert into a database using `psycopg2.extras.execute_values`, iterate directly over the list of dictionaries instead of converting it to a Pandas DataFrame and using `iterrows()`. This saves memory and is significantly faster.
50 changes: 28 additions & 22 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,34 @@ 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

# ⚡ Bolt Optimization: Iterate directly over list of dicts instead of using pd.DataFrame(data) and iterrows()
# Pandas iterrows() is notoriously slow because it converts each row to a Series object.
# Iterating over the list of dicts directly is significantly faster and saves memory.
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 +120,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 +137,41 @@ 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