Saltar al contenido

lección 7

Proyecto: pipeline Python que ingiere, limpia y exporta datos

Construye un mini-pipeline profesional completo: API → limpieza → validación → Parquet.

60 min

### El momento de la verdad: todo junto

Has aprendido Pandas, lectura de formatos, limpieza, conexión a bases de datos, APIs, testing y logging. Ahora es el momento de juntar TODO en un pipeline profesional completo. Este es el tipo de código que escribirás en tu primer día de trabajo como ingeniero de datos: un script que ingesta datos de una fuente externa, los limpia, valida su calidad, y los exporta en formato optimizado.

El proyecto: vamos a construir un pipeline que ingiere datos de la API de JSONPlaceholder (posts y usuarios), los combina, limpia, valida con tests automatizados y exporta como Parquet. Es un microcosmos de lo que harás con datos reales de producción.

### Arquitectura del pipeline

Pipeline profesional: 5 pasos claros con logging en cada uno

### Estructura del proyecto

  • pipeline_posts/ — carpeta raíz
  • pipeline_posts/src/ — código fuente
  • pipeline_posts/src/ingesta.py — funciones de extracción de APIs
  • pipeline_posts/src/limpieza.py — funciones de limpieza y transformación
  • pipeline_posts/src/validacion.py — reglas de calidad de datos
  • pipeline_posts/src/exportar.py — guardar en Parquet/CSV
  • pipeline_posts/src/pipeline.py — orquestador principal
  • pipeline_posts/tests/ — tests automatizados
  • pipeline_posts/tests/test_limpieza.py
  • pipeline_posts/tests/test_validacion.py
  • pipeline_posts/output/ — archivos de salida
  • pipeline_posts/pipeline.log — log de ejecución
  • pipeline_posts/requirements.txt — dependencias

### Paso 1: El módulo de ingesta (src/ingesta.py)

1"""Módulo de ingesta: extrae datos de APIs externas."""
2import requests
3import pandas as pd
4import logging
5import time
6
7logger = logging.getLogger('pipeline.ingesta')
8
9
10def pedir_con_reintentos(url: str, max_intentos: int = 3) -> list:
11 """Petición GET robusta con reintentos y backoff exponencial."""
12 for intento in range(max_intentos):
13 try:
14 response = requests.get(url, timeout=10)
15 response.raise_for_status()
16 return response.json()
17 except (requests.exceptions.RequestException) as e:
18 if intento < max_intentos - 1:
19 espera = 2 ** intento
20 logger.warning(f"Intento {intento+1} fallido: {e}. Esperando {espera}s...")
21 time.sleep(espera)
22 else:
23 logger.error(f"Todos los intentos fallaron para {url}")
24 raise
25
26
27def ingerir_usuarios() -> pd.DataFrame:
28 """Ingiere usuarios de JSONPlaceholder."""
29 logger.info("Ingiriendo usuarios...")
30 data = pedir_con_reintentos('https://jsonplaceholder.typicode.com/users')
31 df = pd.json_normalize(data, sep='_')
32 df = df[['id', 'name', 'email', 'address_city', 'company_name']].rename(columns={
33 'id': 'user_id',
34 'name': 'autor_nombre',
35 'address_city': 'autor_ciudad',
36 'company_name': 'autor_empresa'
37 })
38 logger.info(f"Ingeridos {len(df)} usuarios")
39 return df
40
41
42def ingerir_posts() -> pd.DataFrame:
43 """Ingiere posts de JSONPlaceholder."""
44 logger.info("Ingiriendo posts...")
45 data = pedir_con_reintentos('https://jsonplaceholder.typicode.com/posts')
46 df = pd.DataFrame(data)
47 df = df.rename(columns={'userId': 'user_id', 'id': 'post_id'})
48 logger.info(f"Ingeridos {len(df)} posts")
49 return df

src/ingesta.py: funciones independientes, reutilizables y con logging

### Paso 2: El módulo de limpieza (src/limpieza.py)

1"""Módulo de limpieza: transforma datos crudos en datos usables."""
2import pandas as pd
3import logging
4
5logger = logging.getLogger('pipeline.limpieza')
6
7
8def merge_posts_usuarios(posts: pd.DataFrame, usuarios: pd.DataFrame) -> pd.DataFrame:
9 """Combina posts con información del autor."""
10 logger.info(f"Combinando {len(posts)} posts con {len(usuarios)} usuarios...")
11 df = posts.merge(usuarios, on='user_id', how='left')
12 sin_autor = df['autor_nombre'].isnull().sum()
13 if sin_autor > 0:
14 logger.warning(f"{sin_autor} posts sin autor asociado")
15 return df
16
17
18def limpiar_posts(df: pd.DataFrame) -> pd.DataFrame:
19 """Limpia el DataFrame de posts enriquecido."""
20 df = df.copy() # No modificar el DataFrame original
21 filas_antes = len(df)
22 logger.info(f"Limpiando {filas_antes} posts...")
23
24 # 1. Eliminar duplicados
25 df = df.drop_duplicates(subset=['post_id'], keep='first')
26
27 # 2. Eliminar posts sin título
28 df = df.dropna(subset=['title'])
29
30 # 3. Normalizar texto
31 df['title'] = df['title'].str.strip().str.capitalize()
32 df['autor_nombre'] = df['autor_nombre'].str.strip().str.title()
33 df['autor_ciudad'] = df['autor_ciudad'].fillna('Desconocida').str.strip().str.title()
34
35 # 4. Crear columna de longitud del body (métrica útil)
36 df['body_length'] = df['body'].str.len()
37
38 # 5. Eliminar columna email (PII - no la necesitamos en el output)
39 if 'email' in df.columns:
40 df = df.drop(columns=['email'])
41
42 filas_despues = len(df)
43 logger.info(f"Limpieza completa: {filas_antes}{filas_despues} filas")
44
45 return df.reset_index(drop=True)

src/limpieza.py: merge + transformación con métricas de lo que se eliminó

### Paso 3: El módulo de validación (src/validacion.py)

1"""Módulo de validación: reglas de calidad que los datos DEBEN cumplir."""
2import pandas as pd
3import logging
4
5logger = logging.getLogger('pipeline.validacion')
6
7
8class ValidacionError(Exception):
9 """Los datos no pasan las reglas de calidad."""
10 pass
11
12
13def validar_posts(df: pd.DataFrame) -> None:
14 """Valida el DataFrame de posts contra reglas de calidad."""
15 errores = []
16
17 # Regla 1: No vacío
18 if df.empty:
19 errores.append("DataFrame vacío")
20
21 # Regla 2: Columnas obligatorias
22 columnas_requeridas = {'post_id', 'title', 'body', 'user_id', 'autor_nombre'}
23 faltantes = columnas_requeridas - set(df.columns)
24 if faltantes:
25 errores.append(f"Columnas faltantes: {faltantes}")
26
27 # Si faltan columnas, las reglas siguientes no se pueden evaluar
28 if errores:
29 msg = "Validación fallida:\n" + "\n".join(f" - {e}" for e in errores)
30 logger.error(msg)
31 raise ValidacionError(msg)
32
33 # A partir de aquí las columnas existen
34 # Regla 3: Sin duplicados en post_id
35 duplicados = df['post_id'].duplicated().sum()
36 if duplicados > 0:
37 errores.append(f"{duplicados} post_id duplicados")
38
39 # Regla 4: Títulos no vacíos
40 titulos_vacios = df['title'].str.strip().eq('').sum()
41 if titulos_vacios > 0:
42 errores.append(f"{titulos_vacios} títulos vacíos")
43
44 # Regla 5: body_length razonable (> 0)
45 if 'body_length' in df.columns:
46 bodies_vacios = (df['body_length'] == 0).sum()
47 if bodies_vacios > 0:
48 errores.append(f"{bodies_vacios} posts con body vacío")
49
50 if errores:
51 msg = "Validación fallida:\n" + "\n".join(f" - {e}" for e in errores)
52 logger.error(msg)
53 raise ValidacionError(msg)
54
55 logger.info(f"Validación exitosa: {len(df)} posts pasan todas las reglas")

src/validacion.py: reglas de calidad que actúan como guardián antes de exportar

### Paso 4: El módulo de exportación (src/exportar.py)

1"""Módulo de exportación: escribe el resultado en los formatos de salida."""
2import sqlite3
3from pathlib import Path
4import pandas as pd
5import logging
6
7logger = logging.getLogger('pipeline.exportar')
8
9
10def exportar(df: pd.DataFrame, carpeta: str = 'output') -> None:
11 """Exporta el DataFrame a Parquet, CSV y una tabla SQLite."""
12 destino = Path(carpeta)
13 destino.mkdir(exist_ok=True)
14
15 # Formatos planos
16 df.to_parquet(destino / 'posts_enriquecidos.parquet', index=False)
17 df.to_csv(destino / 'posts_enriquecidos.csv', index=False)
18
19 # Tabla para el equipo de BI (patrón de la lección de bases de datos)
20 with sqlite3.connect(destino / 'posts.db') as conn:
21 df.to_sql('posts_enriquecidos', conn, if_exists='replace', index=False)
22 conn.commit() # toda operación que MODIFICA datos acaba en commit
23
24 logger.info(f"Exportados {len(df)} posts a {destino}/ (parquet, csv, sqlite)")

src/exportar.py: guarda el resultado en tres formatos, incluida una tabla consultable

### Paso 5: El orquestador (src/pipeline.py)

1"""Pipeline principal: orquesta todos los pasos."""
2import logging
3import time
4import os
5from pathlib import Path
6
7logging.basicConfig(
8 level=logging.INFO,
9 format='%(asctime)s | %(levelname)-8s | %(name)-20s | %(message)s',
10 datefmt='%Y-%m-%d %H:%M:%S',
11 handlers=[
12 logging.FileHandler('pipeline.log'),
13 logging.StreamHandler()
14 ]
15)
16logger = logging.getLogger('pipeline')
17
18from ingesta import ingerir_usuarios, ingerir_posts
19from limpieza import merge_posts_usuarios, limpiar_posts
20from validacion import validar_posts
21
22
23def ejecutar():
24 """Ejecuta el pipeline completo."""
25 inicio = time.time()
26 logger.info("=" * 60)
27 logger.info("INICIO: Pipeline de posts enriquecidos")
28 logger.info("=" * 60)
29
30 output_dir = Path('output')
31 output_dir.mkdir(exist_ok=True)
32
33 try:
34 # 1. Ingesta
35 usuarios = ingerir_usuarios()
36 posts = ingerir_posts()
37
38 # 2. Merge
39 df = merge_posts_usuarios(posts, usuarios)
40
41 # 3. Limpieza
42 df = limpiar_posts(df)
43
44 # 4. Validación
45 validar_posts(df)
46
47 # 5. Export
48 parquet_path = output_dir / 'posts_enriquecidos.parquet'
49 csv_path = output_dir / 'posts_enriquecidos.csv'
50
51 df.to_parquet(parquet_path, index=False)
52 df.to_csv(csv_path, index=False)
53
54 duracion = time.time() - inicio
55 logger.info(f"Exportados {len(df)} posts a {parquet_path}")
56 logger.info(f"ÉXITO: Pipeline completado en {duracion:.1f}s")
57 logger.info("=" * 60)
58
59 return df
60
61 except Exception as e:
62 duracion = time.time() - inicio
63 logger.critical(f"FALLO tras {duracion:.1f}s: {type(e).__name__}: {e}", exc_info=True)
64 raise
65
66
67if __name__ == '__main__':
68 ejecutar()

src/pipeline.py: el director de orquesta que llama a cada módulo en orden

Observa la separación: cada módulo hace UNA cosa (ingesta, limpieza, validación, export). El pipeline solo ORQUESTA. Si mañana cambias de API, tocas solo ingesta.py. Si cambias reglas de calidad, tocas solo validacion.py. Esto es el principio de responsabilidad única aplicado a pipelines de datos.

### Paso 5: Los tests (tests/test_limpieza.py)

1"""Tests para el módulo de limpieza."""
2import pandas as pd
3import pytest
4from src.limpieza import limpiar_posts, merge_posts_usuarios
5
6
7@pytest.fixture
8def df_posts_raw():
9 return pd.DataFrame({
10 'post_id': [1, 2, 3],
11 'user_id': [1, 1, 2],
12 'title': [' hello world ', 'SECOND POST', 'third'],
13 'body': ['Content one', 'Content two', 'Content three'],
14 'autor_nombre': [' john doe ', ' john doe ', 'jane smith'],
15 'autor_ciudad': [None, 'New York', ' boston '],
16 'autor_empresa': ['Corp A', 'Corp A', 'Corp B']
17 })
18
19
20def test_limpiar_normaliza_titulos(df_posts_raw):
21 resultado = limpiar_posts(df_posts_raw)
22 assert resultado['title'].tolist() == ['Hello world', 'Second post', 'Third']
23
24
25def test_limpiar_ciudad_nula_es_desconocida(df_posts_raw):
26 resultado = limpiar_posts(df_posts_raw)
27 assert 'Desconocida' in resultado['autor_ciudad'].values
28
29
30def test_limpiar_crea_body_length(df_posts_raw):
31 resultado = limpiar_posts(df_posts_raw)
32 assert 'body_length' in resultado.columns
33 assert (resultado['body_length'] > 0).all()
34
35
36def test_limpiar_sin_duplicados():
37 df = pd.DataFrame({
38 'post_id': [1, 1, 2],
39 'user_id': [1, 1, 2],
40 'title': ['A', 'A duplicado', 'B'],
41 'body': ['x', 'x', 'y'],
42 'autor_nombre': ['Ana', 'Ana', 'Carlos'],
43 'autor_ciudad': ['Madrid', 'Madrid', 'Barcelona'],
44 'autor_empresa': ['X', 'X', 'Y']
45 })
46 resultado = limpiar_posts(df)
47 assert len(resultado) == 2

Tests que verifican cada regla de limpieza de forma aislada

### Ejecutar el proyecto

1# Instalar dependencias
2pip install -r requirements.txt
3
4# Ejecutar tests PRIMERO (siempre)
5pytest tests/ -v
6
7# Si los tests pasan, ejecutar el pipeline
8python src/pipeline.py
9
10# Ver los logs
11cat pipeline.log
12
13# Verificar el output
14python -c "import pandas as pd; df = pd.read_parquet('output/posts_enriquecidos.parquet'); print(df.info())"

Flujo de ejecución: tests → pipeline → verificar output

### Lo que has construido: visión completa

Mira lo que has creado: un pipeline profesional con ingesta robusta (reintentos, timeout), transformación modular (merge, limpieza), validación de calidad (reglas de negocio), export multi-formato (Parquet + CSV), logging completo (cada paso, tiempos, errores) y tests automatizados (verificación antes de desplegar). Esto es exactamente lo que se espera de un ingeniero de datos junior en su primera semana.

  1. 01.Modularidad: cada archivo hace una cosa. Fácil de mantener, testear y extender.
  2. 02.Robustez: reintentos en la ingesta, manejo de nulos, validación de calidad.
  3. 03.Observabilidad: logs que te dicen qué pasó sin mirar el código.
  4. 04.Reproducibilidad: requirements.txt fija versiones. El mismo resultado cada vez.
  5. 05.Calidad: tests que verifican que no rompes nada al cambiar código.

Este pipeline se ejecuta manualmente. En el mundo real, necesitas que corra AUTOMÁTICAMENTE cada día sin que tú estés. Eso es la orquestación (Skill 15): Step Functions o Airflow ejecutan tu pipeline según un schedule, manejan fallos y te alertan. Pero la lógica del pipeline es EXACTAMENTE lo que acabas de escribir.

Si puedes explicar tu pipeline en una servilleta (5 cajas con flechas), está bien diseñado. Si necesitas un diagrama de A3, es demasiado complejo. Simplifica. Los mejores pipelines son aburridos: hacen una cosa simple muy bien, cada día, sin sorpresas.

## ejercicios

[01]

Pipeline completo: datos del clima de España

Construye un mini-pipeline que: 1) Ingiera datos del tiempo de Open-Meteo para 5 ciudades españolas, 2) Combine todo en un DataFrame con columna "ciudad", 3) Limpie y valide (no nulos en temperatura), 4) Exporte a Parquet. Incluye logging y reintentos.

Cargando editor...
[02]

Escribir tests de validación para el pipeline de posts

Escribe 4 tests para la función validar_posts(): test con DataFrame vacío (debe lanzar ValidacionError), test con columna faltante, test con duplicados en post_id, y test que pasa correctamente con datos buenos.

Cargando editor...
[03]

Crear el requirements.txt del proyecto

Lista todas las dependencias directas del proyecto pipeline_posts con versiones fijadas. Incluye: pandas, requests, pytest y pyarrow (para Parquet). Solo las que tu código importa.

Cargando editor...

Regístrate para guardar tu progreso.

## comentarios

Reporta erratas, ayuda a otros o comparte tu opinión. Sé constructivo.

Inicia sesión para comentar y responder.

cargando comentarios...