From 7901a8a41fe13c8383e5b7cbe123c399a5c2a3d7 Mon Sep 17 00:00:00 2001 From: "google-labs-jules[bot]" <161369871+google-labs-jules[bot]@users.noreply.github.com> Date: Thu, 16 Jul 2026 05:00:37 +0000 Subject: [PATCH] =?UTF-8?q?=E2=9A=A1=20Bolt:=20Pandas=20iterrows()=20optim?= =?UTF-8?q?ization=20in=20ETL=20pipeline?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Optimizes the database loading step in the public data ETL pipeline by bypassing Pandas DataFrame creation and iterating directly over the raw list of dictionaries. This avoids the significant overhead of Pandas iterrows(). Co-authored-by: Vagarh <111590756+Vagarh@users.noreply.github.com> --- .jules/bolt.md | 3 ++ .../dags/public_data_etl.py | 50 +++++++++++-------- 2 files changed, 31 insertions(+), 22 deletions(-) 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