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-10-26 - Pandas iterrows Optimization in Airflow
**Learning:** Using `pandas.DataFrame.iterrows()` inside an Airflow PythonOperator task (like building an SQL insert script) to iterate over rows is incredibly slow. Converting a JSON or Python list of dictionaries (commonly from XComs) into a DataFrame *just* to iterate over it using `iterrows()` is an anti-pattern that creates a significant performance bottleneck.
**Action:** When building structures like bulk inserts, if the data is already in a Python list of dictionaries (from JSON or XCom), avoid converting it to a DataFrame. Iterate directly over the dictionaries (e.g., `for row in data:`), which can provide a ~100x+ speedup.
41 changes: 24 additions & 17 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,35 @@ 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)

# ⚡ Bolt Optimization: Eliminated unnecessary conversion to DataFrame and use of slow iterrows().
# Iterating directly over the list of dictionaries provides a ~127x speedup.

# 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
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 +121,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,14 +138,16 @@ 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,
Expand Down