lección 4
Tu primer workflow: ingesta → transformación → carga
Construye un pipeline real con Step Functions que ingiere datos de S3, los transforma y los carga en destino.
⏱ 60 min
### De Pass a Task: conectando código real
En la lección anterior prototipamos un flujo con estados Pass — simulaciones sin código real. Ahora vamos a dar el salto: crear estados Task que ejecutan funciones Lambda reales. Es el momento donde tu máquina de estado deja de ser un diagrama bonito y empieza a HACER cosas.
El escenario: una empresa de e-commerce recibe cada día un archivo CSV con las ventas del día anterior en un bucket S3. Necesitamos un pipeline que automáticamente: 1) Detecte que el archivo llegó, 2) Lo lea y transforme (limpiar, calcular métricas), 3) Guarde el resultado transformado en otro bucket (zona gold). Todo orquestado por Step Functions.
### Paso 1: Preparar la infraestructura en LocalStack
Primero necesitamos crear los buckets S3 y subir un archivo de prueba. Recuerda: todo en LocalStack, sin coste, sin cuenta AWS.
1# Crear buckets S3 en LocalStack2aws s3 mb s3://datos-raw --endpoint-url http://localhost:45663aws s3 mb s3://datos-gold --endpoint-url http://localhost:456645# Crear un CSV de prueba6cat > ventas_20240115.csv << 'EOF'7fecha,producto,cantidad,precio_unitario,cliente_id82024-01-15,Laptop Pro,2,1299.99,C00192024-01-15,Mouse Wireless,5,29.99,C002102024-01-15,Teclado Mecánico,3,89.99,C001112024-01-15,Monitor 4K,1,499.99,C003122024-01-15,Cable USB-C,,12.99,C002132024-01-15,Webcam HD,2,79.99,14EOF1516# Subir a S3 raw17aws s3 cp ventas_20240115.csv s3://datos-raw/ventas/2024/01/15/ \18 --endpoint-url http://localhost:45661920# Verificar21aws s3 ls s3://datos-raw/ventas/2024/01/15/ --endpoint-url http://localhost:4566
Nota: el CSV tiene datos sucios a propósito (cantidad vacía, cliente_id vacío)
### Paso 2: Crear las funciones Lambda
Cada paso del pipeline será una función Lambda independiente. Esto sigue el principio de responsabilidad única: cada función hace UNA cosa. Step Functions las coordina.
1# lambda_ingesta.py — Lee el CSV de S3 y valida que tiene datos2import json3import os4import boto35import csv6import io78ENDPOINT = os.environ.get('AWS_ENDPOINT_URL', 'http://localhost:4566')9s3 = boto3.client('s3', endpoint_url=ENDPOINT)1011def handler(event, context):12 """Lee un CSV de S3 y devuelve los datos como lista de dicts."""13 bucket = event['bucket']14 key = event['key']1516 # Leer archivo de S317 response = s3.get_object(Bucket=bucket, Key=key)18 content = response['Body'].read().decode('utf-8')1920 # Parsear CSV21 reader = csv.DictReader(io.StringIO(content))22 registros = list(reader)2324 if not registros:25 raise Exception(f"Archivo vacío: s3://{bucket}/{key}")2627 return {28 "status": "ok",29 "registros": len(registros),30 "datos": registros,31 "fuente": f"s3://{bucket}/{key}"32 }
Lambda de ingesta: lee, parsea y valida. Si está vacío, falla (y Step Functions gestiona el error)
1# lambda_transformar.py — Limpia y calcula métricas2import json34def handler(event, context):5 """Limpia datos y calcula el total por registro."""6 datos = event['datos']78 limpios = []9 descartados = 01011 for registro in datos:12 # Descartar registros sin cantidad o sin cliente13 if not registro.get('cantidad') or not registro.get('cliente_id'):14 descartados += 115 continue1617 # Calcular total18 registro['precio_unitario'] = float(registro['precio_unitario'])19 registro['total'] = round(registro['precio_unitario'] * float(registro['cantidad']), 2)20 registro['cantidad'] = int(registro['cantidad'])21 limpios.append(registro)2223 return {24 "status": "ok",25 "registros_limpios": len(limpios),26 "registros_descartados": descartados,27 "datos_transformados": limpios28 }
Lambda de transformación: limpia registros inválidos y calcula totales
1# lambda_cargar.py — Escribe el resultado en S3 Gold2import json3import os4import boto356ENDPOINT = os.environ.get('AWS_ENDPOINT_URL', 'http://localhost:4566')7s3 = boto3.client('s3', endpoint_url=ENDPOINT)89def handler(event, context):10 """Escribe los datos transformados en S3 gold como JSON."""11 datos = event['datos_transformados']12 fecha = event.get('fecha', 'unknown')1314 # Escribir en S3 Gold (JSON lines para fácil lectura)15 output = ''.join(json.dumps(r, ensure_ascii=False) + '\n' for r in datos)1617 key = f"ventas/fecha={fecha}/datos.jsonl"18 s3.put_object(19 Bucket='datos-gold',20 Key=key,21 Body=output.encode('utf-8')22 )2324 return {25 "status": "ok",26 "destino": f"s3://datos-gold/{key}",27 "registros_cargados": len(datos)28 }
Lambda de carga: escribe en S3 Gold particionado por fecha
### Paso 3: Desplegar las Lambdas en LocalStack
1# Empaquetar cada Lambda como ZIP2zip lambda_ingesta.zip lambda_ingesta.py3zip lambda_transformar.zip lambda_transformar.py4zip lambda_cargar.zip lambda_cargar.py56# Crear las funciones en LocalStack7aws lambda create-function \8 --function-name "ingesta" \9 --runtime "python3.11" \10 --handler "lambda_ingesta.handler" \11 --zip-file fileb://lambda_ingesta.zip \12 --role "arn:aws:iam::000000000000:role/DummyRole" \13 --endpoint-url http://localhost:45661415aws lambda create-function \16 --function-name "transformar" \17 --runtime "python3.11" \18 --handler "lambda_transformar.handler" \19 --zip-file fileb://lambda_transformar.zip \20 --role "arn:aws:iam::000000000000:role/DummyRole" \21 --endpoint-url http://localhost:45662223aws lambda create-function \24 --function-name "cargar" \25 --runtime "python3.11" \26 --handler "lambda_cargar.handler" \27 --zip-file fileb://lambda_cargar.zip \28 --role "arn:aws:iam::000000000000:role/DummyRole" \29 --endpoint-url http://localhost:4566
Tres funciones Lambda independientes, cada una con su responsabilidad
### Paso 4: La máquina de estado que lo orquesta todo
Ahora viene la parte bonita: definir la máquina de estado que coordina las tres Lambdas. Fíjate en cómo cada Task referencia una Lambda por su ARN y usa ResultPath para acumular resultados.
1{2 "Comment": "Pipeline ETL: Ingesta → Transformación → Carga",3 "StartAt": "Ingesta",4 "States": {5 "Ingesta": {6 "Type": "Task",7 "Resource": "arn:aws:lambda:eu-west-1:000000000000:function:ingesta",8 "Parameters": {9 "bucket": "datos-raw",10 "key.$": "$.key"11 },12 "ResultPath": "$.ingesta_resultado",13 "Next": "Transformacion"14 },15 "Transformacion": {16 "Type": "Task",17 "Resource": "arn:aws:lambda:eu-west-1:000000000000:function:transformar",18 "Parameters": {19 "datos.$": "$.ingesta_resultado.datos"20 },21 "ResultPath": "$.transformacion_resultado",22 "Next": "Carga"23 },24 "Carga": {25 "Type": "Task",26 "Resource": "arn:aws:lambda:eu-west-1:000000000000:function:cargar",27 "Parameters": {28 "datos_transformados.$": "$.transformacion_resultado.datos_transformados",29 "fecha.$": "$.fecha"30 },31 "ResultPath": "$.carga_resultado",32 "End": true33 }34 }35}
La notación .$ indica "toma el valor de este JsonPath del input". Es como pasar parámetros entre funciones.
Fíjate en el patrón: Parameters filtra QUÉ recibe cada Lambda (no le pasa todo el estado, solo lo que necesita). ResultPath decide DÓNDE se guarda el resultado. Este patrón es el pan de cada día con Step Functions.
### Paso 5: Crear y ejecutar el pipeline completo
1# Crear la máquina de estado2aws stepfunctions create-state-machine \3 --name "PipelineVentasETL" \4 --definition file://pipeline-etl.json \5 --role-arn "arn:aws:iam::000000000000:role/DummyRole" \6 --endpoint-url http://localhost:456678# Ejecutar con los datos del 15 de enero9aws stepfunctions start-execution \10 --state-machine-arn "arn:aws:states:eu-west-1:000000000000:stateMachine:PipelineVentasETL" \11 --input '{"key": "ventas/2024/01/15/ventas_20240115.csv", "fecha": "2024-01-15"}' \12 --endpoint-url http://localhost:45661314# Verificar que los datos llegaron a Gold15aws s3 ls s3://datos-gold/ventas/ --recursive --endpoint-url http://localhost:4566
El pipeline completo en acción: de CSV sucio en raw a JSON limpio en gold
Acabas de construir tu primer pipeline orquestado de verdad. No es un script monolítico que depende de que lo ejecutes tú. Es un sistema declarativo donde el flujo está separado del código, cada paso es independiente, y Step Functions se encarga de la coordinación. Si mañana necesitas añadir un paso de "notificar por email", solo añades un estado más al JSON sin tocar las Lambdas existentes.
Ojo con el tamaño de los datos entre estados. Aqui pasamos los registros dentro del estado porque son seis y se ve mejor. Pero el input o output de un estado en Step Functions no puede pasar de 256 KiB (es un limite duro, no se puede subir). Con la forma de fila de este CSV, eso son unas mil filas — un dia de ventas de una tienda pequeña. En cuanto tu pipeline procese datos de verdad, lo que viaja entre estados es la ruta en S3 (la referencia), no los datos. S3 es el bus de datos; Step Functions solo pasa los billetes.
Y si algo falla, no te quedes mirando un bucket vacio sin pistas. Usa describe-execution para ver en que estado se rompio y get-execution-history para ver que input recibio ese estado. Y en la proxima leccion veras Retry y Catch: la respuesta de Step Functions a los fallos transitorios.
En producción real (no LocalStack), las Lambdas tienen un límite de 15 minutos de ejecución y 10GB de memoria. Si tu transformación tarda más o necesita más RAM, necesitas ECS/Fargate o Glue en lugar de Lambda. Step Functions soporta ambos — solo cambia el Resource del Task.
## ejercicios
Añadir un paso Wait antes de la carga
El equipo de DBA te pide que esperes 30 segundos entre la transformación y la carga para no saturar la base de datos. Modifica la máquina de estado para añadir un estado Wait entre Transformacion y Carga.
Escribir una Lambda de transformación
Escribe una función Lambda que reciba una lista de registros de ventas y calcule: total por cliente, número de productos distintos por cliente, y descarte registros con precio_unitario negativo o campos vacíos.
💡 Resultado esperado
{
"status": "ok",
"metricas_por_cliente": [
{Verificar los datos en S3 Gold
Después de ejecutar el pipeline, escribe los comandos para: 1) listar archivos en s3://datos-gold, 2) descargar el resultado, 3) verificar que contiene los registros esperados.
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...