Saltar al contenido

lección 13

Proyecto: pipeline end-to-end en AWS

Construye paso a paso un pipeline completo con LocalStack: S3 como lake, Lambda para ingesta, Glue para catálogo, Athena para consultas y Step Functions para orquestación.

65 min

Has llegado al final del track. Esta lección es diferente: no hay teoría nueva. Todo lo que necesitas ya lo sabes. Aquí vas a CONSTRUIR un pipeline de datos completo end-to-end usando todos los servicios que has aprendido, corriendo íntegramente en LocalStack en tu portátil.

El escenario: FashionStore recibe archivos CSV de ventas cada hora desde su sistema de punto de venta. Necesitan un pipeline que: 1) detecte automáticamente los archivos nuevos, 2) los valide, 3) los transforme a Parquet, 4) los registre en el catálogo, 5) permita consultas SQL, y 6) notifique si algo falla. Todo serverless, todo automatizado, todo con los servicios que has aprendido.

### Arquitectura del proyecto

El pipeline completo que vas a construir en esta lección

### Paso 1: Levantar LocalStack y crear la infraestructura base

1import boto3
2import json
3
4# Helper para crear clientes
5def client(service):
6 return boto3.client(service, endpoint_url='http://localhost:4566',
7 aws_access_key_id='test', aws_secret_access_key='test',
8 region_name='eu-west-1')
9
10s3 = client('s3')
11iam = client('iam')
12glue = client('glue')
13sfn = client('stepfunctions')
14lambda_client = client('lambda')
15
16# Crear buckets
17for bucket in ['fashionstore-raw', 'fashionstore-processed', 'fashionstore-athena-results']:
18 s3.create_bucket(Bucket=bucket,
19 CreateBucketConfiguration={'LocationConstraint': 'eu-west-1'})
20 print(f"✅ Bucket: {bucket}")
21
22# Crear database en Glue Catalog
23glue.create_database(DatabaseInput={'Name': 'fashionstore_prod'})
24print("✅ Database: fashionstore_prod")

Paso 1: infraestructura base — buckets + database del catálogo

### Paso 2: Crear roles IAM con mínimo privilegio

1# Role para Lambdas del pipeline
2trust = json.dumps({"Version": "2012-10-17", "Statement": [
3 {"Effect": "Allow", "Principal": {"Service": "lambda.amazonaws.com"}, "Action": "sts:AssumeRole"}
4]})
5
6iam.create_role(RoleName='pipeline-lambda-role', AssumeRolePolicyDocument=trust)
7
8# Policy: las Lambdas pueden leer raw, escribir processed, y actualizar catálogo
9policy_doc = json.dumps({"Version": "2012-10-17", "Statement": [
10 {"Effect": "Allow", "Action": ["s3:GetObject"], "Resource": ["arn:aws:s3:::fashionstore-raw/*"]},
11 {"Effect": "Allow", "Action": ["s3:PutObject"], "Resource": ["arn:aws:s3:::fashionstore-processed/*"]},
12 {"Effect": "Allow", "Action": ["glue:*"], "Resource": ["*"]},
13 {"Effect": "Allow", "Action": ["logs:*"], "Resource": ["*"]},
14]})
15
16policy = iam.create_policy(PolicyName='pipeline-lambda-policy', PolicyDocument=policy_doc)
17iam.attach_role_policy(RoleName='pipeline-lambda-role', PolicyArn=policy['Policy']['Arn'])
18print("✅ IAM configurado con mínimo privilegio")

Paso 2: roles con permisos mínimos — cada Lambda solo accede a lo que necesita

### Paso 3: Crear las funciones Lambda

1import zipfile, io
2
3def crear_lambda(nombre, codigo):
4 """Helper para crear una Lambda con su código."""
5 buf = io.BytesIO()
6 with zipfile.ZipFile(buf, 'w') as zf:
7 zf.writestr('handler.py', codigo)
8 buf.seek(0)
9 lambda_client.create_function(
10 FunctionName=nombre, Runtime='python3.11',
11 Role='arn:aws:iam::000000000000:role/pipeline-lambda-role',
12 Handler='handler.handler', Code={'ZipFile': buf.read()},
13 Timeout=300, MemorySize=256
14 )
15 print(f" ✅ Lambda: {nombre}")
16
17# Lambda 1: Validar archivo
18crear_lambda('fs-validar', """
19import json
20def handler(event, context):
21 key = event['key']
22 size = event.get('size', 0)
23 errores = []
24 if size == 0: errores.append('Archivo vacío')
25 if not key.endswith('.csv'): errores.append('No es CSV')
26 return {'key': key, 'valido': len(errores) == 0, 'errores': errores}
27""")
28
29# Lambda 2: Transformar CSV → JSON (simula Parquet)
30crear_lambda('fs-transformar', """
31import json, csv, io, boto3
32def handler(event, context):
33 s3 = boto3.client('s3', endpoint_url='http://localhost:4566')
34 key = event['key']
35 obj = s3.get_object(Bucket='fashionstore-raw', Key=key)
36 content = obj['Body'].read().decode()
37 reader = csv.DictReader(io.StringIO(content))
38 records = list(reader)
39 out_key = key.replace('raw/', '').replace('.csv', '.json')
40 s3.put_object(Bucket='fashionstore-processed', Key=out_key,
41 Body=json.dumps(records).encode())
42 return {'output_key': out_key, 'records': len(records)}
43""")
44
45# Lambda 3: Registrar en catálogo
46crear_lambda('fs-catalogar', """
47import json, boto3
48def handler(event, context):
49 return {'catalogado': True, 'tabla': 'ventas', 'registros': event.get('records', 0)}
50""")
51
52# Lambda 4: Notificar resultado
53crear_lambda('fs-notificar', """
54import json
55def handler(event, context):
56 print(f"NOTIFICACION: Pipeline completado - {json.dumps(event)}")
57 return {'notificado': True}
58""")
59print("\n✅ Todas las Lambdas creadas")

Paso 3: 4 Lambdas independientes — cada una hace una cosa bien

### Paso 4: Crear la máquina de estados

1# Step Functions: orquesta las 4 Lambdas
2pipeline = {
3 "StartAt": "Validar",
4 "States": {
5 "Validar": {
6 "Type": "Task",
7 "Resource": "arn:aws:lambda:eu-west-1:000000000000:function:fs-validar",
8 "Next": "CheckValidez",
9 "Retry": [{"ErrorEquals": ["States.ALL"], "MaxAttempts": 2}]
10 },
11 "CheckValidez": {
12 "Type": "Choice",
13 "Choices": [{"Variable": "$.valido", "BooleanEquals": True, "Next": "Transformar"}],
14 "Default": "NotificarError"
15 },
16 "Transformar": {
17 "Type": "Task",
18 "Resource": "arn:aws:lambda:eu-west-1:000000000000:function:fs-transformar",
19 "Next": "Catalogar"
20 },
21 "Catalogar": {
22 "Type": "Task",
23 "Resource": "arn:aws:lambda:eu-west-1:000000000000:function:fs-catalogar",
24 "Next": "Notificar"
25 },
26 "Notificar": {
27 "Type": "Task",
28 "Resource": "arn:aws:lambda:eu-west-1:000000000000:function:fs-notificar",
29 "End": True
30 },
31 "NotificarError": {
32 "Type": "Task",
33 "Resource": "arn:aws:lambda:eu-west-1:000000000000:function:fs-notificar",
34 "End": True
35 }
36 }
37}
38
39sfn.create_state_machine(
40 name='fashionstore-etl-pipeline',
41 definition=json.dumps(pipeline),
42 roleArn='arn:aws:iam::000000000000:role/pipeline-lambda-role'
43)
44print("✅ Pipeline Step Functions creado")

Paso 4: Step Functions conecta las Lambdas con lógica de flujo y reintentos

### Paso 5: Simular ingesta y ejecutar el pipeline

1import time
2
3# Simular llegada de un archivo CSV de ventas
4csv_data = """order_id,customer_id,product,quantity,amount,date
5ORD001,CUST42,Camiseta Premium,2,51.98,2024-01-15
6ORD002,CUST17,Zapatillas Runner,1,89.99,2024-01-15
7ORD003,CUST42,Pantalón Slim,1,49.99,2024-01-15
8ORD004,CUST08,Chaqueta Invierno,1,129.99,2024-01-15
9"""
10
11s3.put_object(Bucket='fashionstore-raw', Key='ventas/2024-01-15/ventas_hora_14.csv',
12 Body=csv_data.encode())
13print("📤 Archivo subido a S3 raw")
14
15# Ejecutar el pipeline
16execution = sfn.start_execution(
17 stateMachineArn='arn:aws:states:eu-west-1:000000000000:stateMachine:fashionstore-etl-pipeline',
18 input=json.dumps({
19 "key": "ventas/2024-01-15/ventas_hora_14.csv",
20 "size": len(csv_data),
21 "bucket": "fashionstore-raw"
22 })
23)
24print(f"⚡ Pipeline ejecutándose: {execution['executionArn'].split(':')[-1]}")
25
26# Esperar resultado
27while True:
28 desc = sfn.describe_execution(executionArn=execution['executionArn'])
29 if desc['status'] != 'RUNNING':
30 break
31 time.sleep(1)
32
33print(f"\n{'✅' if desc['status'] == 'SUCCEEDED' else '❌'} Pipeline: {desc['status']}")
34if 'output' in desc:
35 print(f" Output: {desc['output']}")

Paso 5: simular un archivo y ver el pipeline en acción

### Paso 6: Verificar resultados

1# Verificar que los datos llegaron a processed
2processed = s3.list_objects_v2(Bucket='fashionstore-processed')
3print("\n📦 Archivos en fashionstore-processed:")
4for obj in processed.get('Contents', []):
5 print(f" {obj['Key']} ({obj['Size']} bytes)")
6
7# Leer el archivo transformado
8if processed.get('Contents'):
9 key = processed['Contents'][0]['Key']
10 data = s3.get_object(Bucket='fashionstore-processed', Key=key)
11 contenido = json.loads(data['Body'].read().decode())
12 print(f"\n📊 Datos transformados ({len(contenido)} registros):")
13 for registro in contenido[:3]:
14 print(f" {registro}")

Paso 6: verificar que todo el pipeline funcionó correctamente

Consejo de senior: este es el pipeline más simple posible. En producción añadirías: dead letter queues para mensajes fallidos, CloudWatch alarms para métricas, un paso de data quality con validaciones, idempotencia (no reprocesar lo ya procesado), y un dashboard de monitorización. Pero el esqueleto es exactamente este.

### Paso 7: Automatizar con S3 Event → Lambda → Step Functions

El último paso es eliminar la ejecución manual. Configuramos un S3 event notification que detecta archivos nuevos en raw/ y una Lambda "dispatcher" que arranca el Step Functions automáticamente. Así el pipeline se ejecuta solo cada vez que llega un archivo:

1# Lambda dispatcher: recibe evento S3 y arranca Step Functions
2dispatcher_code = """
3import json, boto3
4def handler(event, context):
5 sfn = boto3.client('stepfunctions', endpoint_url='http://localhost:4566')
6 for record in event.get('Records', []):
7 key = record['s3']['object']['key']
8 size = record['s3']['object']['size']
9 sfn.start_execution(
10 stateMachineArn='arn:aws:states:eu-west-1:000000000000:stateMachine:fashionstore-etl-pipeline',
11 input=json.dumps({'key': key, 'size': size, 'bucket': 'fashionstore-raw'})
12 )
13 return {'dispatched': len(event.get('Records', []))}
14"""
15
16crear_lambda('fs-dispatcher', dispatcher_code)
17
18# Configurar S3 para triggear el dispatcher
19s3.put_bucket_notification_configuration(
20 Bucket='fashionstore-raw',
21 NotificationConfiguration={
22 'LambdaFunctionConfigurations': [{
23 'LambdaFunctionArn': 'arn:aws:lambda:eu-west-1:000000000000:function:fs-dispatcher',
24 'Events': ['s3:ObjectCreated:*'],
25 'Filter': {'Key': {'FilterRules': [{'Name': 'suffix', 'Value': '.csv'}]}}
26 }]
27 }
28)
29print("✅ Automatización completa: archivo nuevo → pipeline se ejecuta solo")

El pipeline es ahora 100% automático — zero intervención humana

En este proyecto usamos JSON en vez de Parquet porque LocalStack no ejecuta transformaciones PySpark reales. En producción, la Lambda "transformar" sería un Glue Job que convierte CSV a Parquet. El patrón de orquestación es idéntico.

### ¡Felicidades! Has completado la skill de Cloud AWS

Si has llegado hasta aquí y has construido este pipeline funcionando en tu portátil, dominas el stack de datos en AWS. Has aprendido a almacenar datos (S3), protegerlos (IAM), catalogarlos (Glue), consultarlos (Athena), procesarlos a escala (EMR), construir lógica serverless (Lambda), orquestarla (Step Functions), definir infraestructura como código (Terraform) y controlar costes. Eso es el 90% del día a día de un data engineer que trabaja en la nube de Amazon.

Consejo de senior para tu carrera: el pipeline que acabas de construir es un EXCELENTE proyecto para tu portfolio. Ponlo en GitHub con un README que explique la arquitectura, incluye el diagrama, y documenta cómo ejecutarlo con LocalStack. Cuando un entrevistador te pregunte "¿has trabajado con AWS?", ábrelo y muéstralo funcionando. Eso vale más que cualquier certificación.

## ejercicios

[01]

Crear toda la infraestructura del proyecto

Escribe un script que cree TODA la infraestructura necesaria para el pipeline de FashionStore: buckets, roles, policies, database del catálogo — todo de una sola vez.

Cargando editor...
[02]

Pipeline completo: de CSV a consulta

Implementa el pipeline end-to-end: sube un CSV a raw, ejecuta la validación y transformación, y verifica que los datos están disponibles en processed.

Cargando editor...
[03]

Probar el pipeline con datos inválidos

Verifica que el pipeline maneja correctamente archivos inválidos: vacíos, formato incorrecto, y datos corruptos.

Cargando editor...
[04]

Definir el proyecto en Terraform

Genera el código Terraform completo que crearía toda la infraestructura de este proyecto: buckets, roles, Lambdas y Step Functions.

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