Saltar al contenido

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.

Arquitectura del pipeline: Step Functions coordina tres Lambdas que procesan datos entre buckets S3

### 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 LocalStack
2aws s3 mb s3://datos-raw --endpoint-url http://localhost:4566
3aws s3 mb s3://datos-gold --endpoint-url http://localhost:4566
4
5# Crear un CSV de prueba
6cat > ventas_20240115.csv << 'EOF'
7fecha,producto,cantidad,precio_unitario,cliente_id
82024-01-15,Laptop Pro,2,1299.99,C001
92024-01-15,Mouse Wireless,5,29.99,C002
102024-01-15,Teclado Mecánico,3,89.99,C001
112024-01-15,Monitor 4K,1,499.99,C003
122024-01-15,Cable USB-C,,12.99,C002
132024-01-15,Webcam HD,2,79.99,
14EOF
15
16# Subir a S3 raw
17aws s3 cp ventas_20240115.csv s3://datos-raw/ventas/2024/01/15/ \
18 --endpoint-url http://localhost:4566
19
20# Verificar
21aws 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 datos
2import json
3import os
4import boto3
5import csv
6import io
7
8ENDPOINT = os.environ.get('AWS_ENDPOINT_URL', 'http://localhost:4566')
9s3 = boto3.client('s3', endpoint_url=ENDPOINT)
10
11def 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']
15
16 # Leer archivo de S3
17 response = s3.get_object(Bucket=bucket, Key=key)
18 content = response['Body'].read().decode('utf-8')
19
20 # Parsear CSV
21 reader = csv.DictReader(io.StringIO(content))
22 registros = list(reader)
23
24 if not registros:
25 raise Exception(f"Archivo vacío: s3://{bucket}/{key}")
26
27 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étricas
2import json
3
4def handler(event, context):
5 """Limpia datos y calcula el total por registro."""
6 datos = event['datos']
7
8 limpios = []
9 descartados = 0
10
11 for registro in datos:
12 # Descartar registros sin cantidad o sin cliente
13 if not registro.get('cantidad') or not registro.get('cliente_id'):
14 descartados += 1
15 continue
16
17 # Calcular total
18 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)
22
23 return {
24 "status": "ok",
25 "registros_limpios": len(limpios),
26 "registros_descartados": descartados,
27 "datos_transformados": limpios
28 }

Lambda de transformación: limpia registros inválidos y calcula totales

1# lambda_cargar.py — Escribe el resultado en S3 Gold
2import json
3import os
4import boto3
5
6ENDPOINT = os.environ.get('AWS_ENDPOINT_URL', 'http://localhost:4566')
7s3 = boto3.client('s3', endpoint_url=ENDPOINT)
8
9def handler(event, context):
10 """Escribe los datos transformados en S3 gold como JSON."""
11 datos = event['datos_transformados']
12 fecha = event.get('fecha', 'unknown')
13
14 # Escribir en S3 Gold (JSON lines para fácil lectura)
15 output = ''.join(json.dumps(r, ensure_ascii=False) + '\n' for r in datos)
16
17 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 )
23
24 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 ZIP
2zip lambda_ingesta.zip lambda_ingesta.py
3zip lambda_transformar.zip lambda_transformar.py
4zip lambda_cargar.zip lambda_cargar.py
5
6# Crear las funciones en LocalStack
7aws 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:4566
14
15aws 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:4566
22
23aws 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.

Step Functions · máquina de estados
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": true
33 }
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 estado
2aws 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:4566
7
8# Ejecutar con los datos del 15 de enero
9aws 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:4566
13
14# Verificar que los datos llegaron a Gold
15aws 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

[01]

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.

Cargando editor...
[02]

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": [
  {
Cargando editor...
[03]

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.

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...