diff --git a/.jules/bolt.md b/.jules/bolt.md index 39e2abf..b1bcfd3 100644 --- a/.jules/bolt.md +++ b/.jules/bolt.md @@ -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. diff --git a/e2e_open_data_pipeline/dags/public_data_etl.py b/e2e_open_data_pipeline/dags/public_data_etl.py index 9595966..0497c4d 100644 --- a/e2e_open_data_pipeline/dags/public_data_etl.py +++ b/e2e_open_data_pipeline/dags/public_data_etl.py @@ -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" @@ -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 = { @@ -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') @@ -77,22 +80,21 @@ 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, @@ -100,10 +102,12 @@ def load_data(**kwargs): ) 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'), @@ -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() @@ -133,7 +137,8 @@ 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 @@ -141,11 +146,12 @@ def load_data(**kwargs): 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: @@ -153,19 +159,19 @@ def load_data(**kwargs): 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