Saltar al contenido

lección 10

Step Functions + Lambda: pipelines serverless completos

Orquestar Lambdas con Step Functions: estados, transiciones, manejo de errores, paralelismo y el patrón de pipeline serverless completo.

⏱ 50 min

Ya conoces Step Functions de la Skill 15 (Orquestación). Ahora lo combinamos con Lambda para construir pipelines de datos completamente serverless. La idea es poderosa: cada paso del pipeline es una Lambda independiente, y Step Functions las coordina — maneja el orden de ejecución, los reintentos si algo falla, las bifurcaciones condicionales, y el paralelismo.

Este es el patrón más moderno y económico para pipelines de datos que no son extremadamente pesados. Si tu pipeline procesa archivos de < 1GB por paso y no necesita Spark, Step Functions + Lambda es más barato y fácil de mantener que Airflow o EMR.

## El patrón de pipeline serverless

  1. 01.S3 Event: llega un archivo nuevo a raw/
  2. 02.Lambda Validadora: verifica formato, tamaño, esquema
  3. 03.Lambda Transformadora: limpia, normaliza, convierte a Parquet
  4. 04.Lambda Cargadora: registra en Glue Catalog o carga en Redshift
  5. 05.Lambda Notificadora: envía Slack/email si todo OK o si falla

## Definir una máquina de estados (ASL)

1import boto3
2import json
3
4sfn = boto3.client('stepfunctions', endpoint_url='http://localhost:4566',
5 aws_access_key_id='test', aws_secret_access_key='test',
6 region_name='eu-west-1')
7
8# Definición de la máquina de estados (ASL - Amazon States Language)
9pipeline_definition = {
10 "Comment": "Pipeline ETL serverless: validar → transformar → cargar",
11 "StartAt": "ValidarArchivo",
12 "States": {
13 "ValidarArchivo": {
14 "Type": "Task",
15 "Resource": "arn:aws:lambda:eu-west-1:000000000000:function:validar-archivo-s3",
16 "Next": "EsValido",
17 "Catch": [{
18 "ErrorEquals": ["States.ALL"],
19 "Next": "NotificarError"
20 }]
21 },
22 "EsValido": {
23 "Type": "Choice",
24 "Choices": [{
25 "Variable": "$.valido",
26 "BooleanEquals": True,
27 "Next": "TransformarDatos"
28 }],
29 "Default": "NotificarError"
30 },
31 "TransformarDatos": {
32 "Type": "Task",
33 "Resource": "arn:aws:lambda:eu-west-1:000000000000:function:transformar-csv",
34 "Next": "CargarEnCatalogo",
35 "Retry": [{
36 "ErrorEquals": ["States.TaskFailed"],
37 "MaxAttempts": 2,
38 "IntervalSeconds": 30
39 }]
40 },
41 "CargarEnCatalogo": {
42 "Type": "Task",
43 "Resource": "arn:aws:lambda:eu-west-1:000000000000:function:registrar-catalogo",
44 "Next": "NotificarExito"
45 },
46 "NotificarExito": {
47 "Type": "Task",
48 "Resource": "arn:aws:lambda:eu-west-1:000000000000:function:notificar",
49 "End": True
50 },
51 "NotificarError": {
52 "Type": "Task",
53 "Resource": "arn:aws:lambda:eu-west-1:000000000000:function:notificar-error",
54 "End": True
55 }
56 }
57}
58
59# Crear la máquina de estados
60sfn.create_state_machine(
61 name='pipeline-etl-serverless',
62 definition=json.dumps(pipeline_definition),
63 roleArn='arn:aws:iam::000000000000:role/StepFunctionsRole'
64)
65print("[OK] Pipeline serverless creado")

Un pipeline completo definido como máquina de estados con reintentos y manejo de errores

## Ejecutar el pipeline y monitorizar

1import time
2
3# Ejecutar el pipeline con un input (simula evento S3)
4execution = sfn.start_execution(
5 stateMachineArn='arn:aws:states:eu-west-1:000000000000:stateMachine:pipeline-etl-serverless',
6 input=json.dumps({
7 "bucket": "fashionstore-datalake",
8 "key": "raw/ventas/2024-01-15.csv",
9 "size": 15000
10 })
11)
12
13execution_arn = execution['executionArn']
14print(f"Ejecutando: {execution_arn}")
15
16# Poll estado
17while True:
18 desc = sfn.describe_execution(executionArn=execution_arn)
19 status = desc['status']
20 print(f" Estado: {status}")
21 if status in ('SUCCEEDED', 'FAILED', 'TIMED_OUT', 'ABORTED'):
22 break
23 time.sleep(2)
24
25if status == 'SUCCEEDED':
26 print(f"[OK] Pipeline completado. Output: {desc.get('output', '')}")
27else:
28 print(f"[ERROR] Pipeline falló: {desc.get('error', 'Unknown')}")

Ejecutar y monitorizar una ejecución del pipeline

## Paralelismo: procesar varios archivos simultáneamente

Step Functions soporta el estado "Parallel" para ejecutar ramas en paralelo, y "Map" para iterar sobre una lista procesando cada elemento concurrentemente. Esto es ideal cuando llegan 10 archivos a la vez: cada uno se procesa en paralelo por su propia Lambda.

1# Estado Map: procesar cada archivo en paralelo
2map_state = {
3 "Type": "Map",
4 "ItemsPath": "$.archivos",
5 "MaxConcurrency": 5, # Máximo 5 en paralelo
6 "Iterator": {
7 "StartAt": "ProcesarArchivo",
8 "States": {
9 "ProcesarArchivo": {
10 "Type": "Task",
11 "Resource": "arn:aws:lambda:...:function:procesar-archivo",
12 "End": True
13 }
14 }
15 },
16 "Next": "ConsolidarResultados"
17}
18print("Map state: procesa N archivos en paralelo (máx 5 concurrentes)")

El estado Map itera sobre una lista procesando items en paralelo

Consejo de senior: Step Functions Express Workflows ($0.000001 por transición) son 10x más baratos que Standard Workflows ($0.000025 por transición) para pipelines de alta frecuencia. Usa Express para pipelines que se ejecutan > 100 veces/día y duran < 5 minutos. Standard para pipelines diarios largos.

Cuidado con los límites de concurrencia de Lambda: por defecto 1000 ejecuciones simultáneas por región. Si tu Step Functions Map lanza 500 Lambdas en paralelo y otros servicios también usan Lambda, puedes alcanzar el límite y provocar throttling. Usa MaxConcurrency en Map para controlarlo.

Cada Lambda es independiente y testeable — Step Functions las conecta

Regístrate para guardar tu progreso.