Módulo 5: AWS Services for AI (S3, Lambda, SageMaker Basics)
4. Integración S3 + Lambda
Descripción
En esta cápsula vas a conectar S3 y Lambda en un flujo event-driven completo: S3 trigger → Lambda → procesa → S3. Hasta ahora trabajaste con S3 (cápsula 02) y Lambda (cápsula 03) por separado. Aquí los combinas. Lambda lee un prompt template de S3, toma un documento nuevo que activó el trigger, construye un prompt contextualizado, invoca el LLM, y escribe la respuesta de vuelta a S3. Es el patrón fundamental de un servicio AI en AWS.
Contexto: En el Módulo 4 construiste un pipeline en LocalStack que integraba S3 y Lambda. Ahora profundizas en los patrones de integración reales: S3 Event Notifications, el formato del evento que recibe Lambda, cómo evitar loops infinitos (Lambda escribe a S3 → S3 triggerear Lambda otra vez), y cómo diseñar flujos que escalen. Al terminar, tendrás un pipeline event-driven que procesa documentos automáticamente.
S3 Event Notifications → Lambda
Cómo funciona el trigger
Cuando un objeto se crea (o se modifica, elimina) en S3, puedes configurar que S3 envíe un evento a Lambda. Lambda se ejecuta automáticamente con la información del objeto:
┌──────────────────────┐
Upload doc ──────→│ S3 Bucket │
PUT documents/ │ ai-assets-xxx/ │
│ documents/new-doc.json│
└──────────┬───────────┘
│ S3 Event Notification
│ (s3:ObjectCreated:*)
↓
┌──────────────────────┐
│ Lambda Function │
│ 1. Lee evento S3 │
│ 2. Lee doc de S3 │
│ 3. Lee prompt de S3 │
│ 4. Invoca LLM │
│ 5. Escribe a S3 │
└──────────┬──────────┘
│ put_object
↓
┌──────────────────────┐
│ S3 Bucket │
│ responses/2026/03/08/│
│ resp-abc123.json │
└──────────────────────┘
Configurar S3 Event Notification con boto3
import boto3
import json
import os
s3 = boto3.client("s3", endpoint_url=os.environ.get("AWS_ENDPOINT_URL"))
lambda_client = boto3.client("lambda", endpoint_url=os.environ.get("AWS_ENDPOINT_URL"))
BUCKET = "ai-assets-dev"
LAMBDA_ARN = os.environ.get(
"PROCESSOR_LAMBDA_ARN",
"arn:aws:lambda:us-east-1:000000000000:function:ai-doc-processor"
)
def configure_s3_trigger(bucket: str, lambda_arn: str, prefix: str) -> dict:
"""Configura S3 para triggerear Lambda cuando se sube un archivo."""
response = s3.put_bucket_notification_configuration(
Bucket=bucket,
NotificationConfiguration={
"LambdaFunctionConfigurations": [
{
"Id": "process-new-documents",
"LambdaFunctionArn": lambda_arn,
"Events": ["s3:ObjectCreated:*"],
"Filter": {
"Key": {
"FilterRules": [
{"Name": "prefix", "Value": prefix},
{"Name": "suffix", "Value": ".json"},
]
}
},
}
]
},
)
print(f"Trigger configurado: {bucket}/{prefix}*.json → {lambda_arn}")
return response
configure_s3_trigger(BUCKET, LAMBDA_ARN, "documents/inbox/")
El evento S3 que recibe Lambda
Cuando S3 triggerear Lambda, el event tiene esta estructura:
# Evento que Lambda recibe cuando se sube un archivo a S3
event = {
"Records": [
{
"eventVersion": "2.1",
"eventSource": "aws:s3",
"awsRegion": "us-east-1",
"eventTime": "2026-03-08T12:00:00.000Z",
"eventName": "ObjectCreated:Put",
"s3": {
"bucket": {
"name": "ai-assets-dev",
"arn": "arn:aws:s3:::ai-assets-dev",
},
"object": {
"key": "documents/inbox/new-report.json",
"size": 4567,
"eTag": "abc123def456",
},
},
}
]
}
Parsear el evento S3 en Lambda
from urllib.parse import unquote_plus
def parse_s3_event(event: dict) -> list[dict]:
"""Extrae información de bucket y key de un evento S3."""
records = []
for record in event.get("Records", []):
s3_info = record.get("s3", {})
bucket = s3_info.get("bucket", {}).get("name", "")
key = unquote_plus(s3_info.get("object", {}).get("key", ""))
size = s3_info.get("object", {}).get("size", 0)
records.append({
"bucket": bucket,
"key": key,
"size": size,
"event_name": record.get("eventName", ""),
"event_time": record.get("eventTime", ""),
"region": record.get("awsRegion", ""),
})
return records
Handler: S3 Trigger → Procesa con LLM → S3
Flujo completo
# s3_trigger_handler.py — Lambda que procesa documentos de S3 con LLM
import json
import logging
import os
import time
from datetime import datetime
from urllib.parse import unquote_plus
import boto3
from openai import OpenAI
logger = logging.getLogger()
logger.setLevel(os.environ.get("LOG_LEVEL", "INFO"))
s3 = boto3.client("s3", endpoint_url=os.environ.get("AWS_ENDPOINT_URL"))
openai_client = OpenAI(api_key=os.environ.get("OPENAI_API_KEY", ""))
BUCKET = os.environ.get("AI_BUCKET", "ai-assets-dev")
PROMPT_KEY = os.environ.get("PROMPT_KEY", "prompts/summarizer/v1/system.txt")
RESPONSE_PREFIX = os.environ.get("RESPONSE_PREFIX", "responses/")
MODEL = os.environ.get("MODEL_NAME", "gpt-4o-mini")
def get_prompt_template() -> str:
"""Lee el prompt template de S3."""
response = s3.get_object(Bucket=BUCKET, Key=PROMPT_KEY)
return response["Body"].read().decode("utf-8")
def get_document(bucket: str, key: str) -> dict:
"""Lee un documento JSON de S3."""
response = s3.get_object(Bucket=bucket, Key=key)
content = response["Body"].read().decode("utf-8")
return json.loads(content)
def save_response(request_id: str, response_data: dict) -> str:
"""Guarda la respuesta de inferencia en S3."""
now = datetime.utcnow()
key = f"{RESPONSE_PREFIX}{now.strftime('%Y/%m/%d')}/{request_id}.json"
s3.put_object(
Bucket=BUCKET,
Key=key,
Body=json.dumps(response_data, ensure_ascii=False).encode("utf-8"),
ContentType="application/json",
)
return key
def process_document(document: dict, prompt_template: str) -> dict:
"""Procesa un documento usando el LLM con el prompt template."""
doc_content = document.get("content", "")
doc_title = document.get("title", "Sin título")
full_prompt = f"{prompt_template}\n\nDocumento: {doc_title}\n\n{doc_content}"
response = openai_client.chat.completions.create(
model=MODEL,
messages=[{"role": "user", "content": full_prompt}],
max_tokens=1000,
)
return {
"answer": response.choices[0].message.content,
"model": response.model,
"tokens_used": response.usage.total_tokens,
"input_tokens": response.usage.prompt_tokens,
"output_tokens": response.usage.completion_tokens,
}
def handler(event, context):
"""Lambda triggerada por S3 que procesa documentos con LLM."""
start_time = time.time()
results = []
prompt_template = get_prompt_template()
for record in event.get("Records", []):
s3_info = record.get("s3", {})
source_bucket = s3_info.get("bucket", {}).get("name", "")
source_key = unquote_plus(s3_info.get("object", {}).get("key", ""))
request_id = f"req-{context.aws_request_id[:8]}"
logger.info(json.dumps({
"event": "processing_document",
"request_id": request_id,
"source": f"s3://{source_bucket}/{source_key}",
}))
try:
document = get_document(source_bucket, source_key)
llm_result = process_document(document, prompt_template)
response_data = {
"request_id": request_id,
"source_key": source_key,
"source_bucket": source_bucket,
"document_title": document.get("title", ""),
"processed_at": datetime.utcnow().isoformat(),
**llm_result,
}
response_key = save_response(request_id, response_data)
logger.info(json.dumps({
"event": "document_processed",
"request_id": request_id,
"response_key": response_key,
"tokens_used": llm_result["tokens_used"],
}))
results.append({
"request_id": request_id,
"source_key": source_key,
"response_key": response_key,
"status": "success",
})
except Exception as e:
logger.error(json.dumps({
"event": "processing_error",
"request_id": request_id,
"source_key": source_key,
"error": str(e),
}))
results.append({
"request_id": request_id,
"source_key": source_key,
"status": "error",
"error": str(e),
})
duration_ms = int((time.time() - start_time) * 1000)
return {
"processed": len(results),
"results": results,
"duration_ms": duration_ms,
}
Evitar Loops Infinitos
El problema
Si Lambda escribe al mismo bucket que la triggerear, puedes crear un loop:
S3 upload (documents/inbox/) → Lambda se ejecuta
→ Lambda escribe a responses/ → ¿S3 triggerear Lambda otra vez?
Si el trigger está configurado para s3:ObjectCreated:* sin filtro de prefijo, Lambda escribe a S3, S3 triggerear Lambda, Lambda escribe a S3... loop infinito. Cada iteración invoca Lambda y cuesta dinero.
Solución 1: Prefijos diferentes (recomendado)
# Trigger SOLO en documents/inbox/ (prefijo de entrada)
# Lambda escribe en responses/ (prefijo de salida)
# El trigger NO se activa en responses/
configure_s3_trigger(
bucket=BUCKET,
lambda_arn=LAMBDA_ARN,
prefix="documents/inbox/", # Solo este prefijo triggerear Lambda
)
Solución 2: Bucket separado para output
INPUT_BUCKET = "ai-input-dev"
OUTPUT_BUCKET = "ai-output-dev"
# Trigger en INPUT_BUCKET
# Lambda escribe en OUTPUT_BUCKET
# Sin posibilidad de loop
Solución 3: Verificar en el handler
def handler(event, context):
for record in event.get("Records", []):
key = record["s3"]["object"]["key"]
if key.startswith("responses/") or key.startswith("processed/"):
logger.info(f"Skipping output file: {key}")
continue
process_document(record)
Flujo Completo: Prompt Template + RAG Document
Pipeline con prompt template dinámico
# pipeline.py — Flujo completo S3 → Lambda → LLM → S3
import boto3
import json
import os
from openai import OpenAI
s3 = boto3.client("s3", endpoint_url=os.environ.get("AWS_ENDPOINT_URL"))
openai_client = OpenAI(api_key=os.environ.get("OPENAI_API_KEY", ""))
BUCKET = os.environ.get("AI_BUCKET", "ai-assets-dev")
def run_ai_pipeline(
prompt_name: str,
prompt_version: str,
document_collection: str,
document_id: str,
max_tokens: int = 1000,
) -> dict:
"""Ejecuta el pipeline completo: lee prompt + documento, invoca LLM, guarda respuesta."""
# 1. Lee prompt template de S3
prompt_key = f"prompts/{prompt_name}/{prompt_version}/system.txt"
prompt_resp = s3.get_object(Bucket=BUCKET, Key=prompt_key)
prompt_template = prompt_resp["Body"].read().decode("utf-8")
# 2. Lee documento de S3
doc_key = f"documents/{document_collection}/{document_id}.json"
doc_resp = s3.get_object(Bucket=BUCKET, Key=doc_key)
document = json.loads(doc_resp["Body"].read().decode("utf-8"))
# 3. Construye prompt contextualizado
doc_content = document.get("content", "")
doc_title = document.get("title", "")
full_prompt = f"{prompt_template}\n\nTítulo: {doc_title}\n\nContenido:\n{doc_content}"
# 4. Invoca LLM
response = openai_client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": full_prompt}],
max_tokens=max_tokens,
)
result = {
"answer": response.choices[0].message.content,
"model": response.model,
"tokens_used": response.usage.total_tokens,
"prompt_name": prompt_name,
"prompt_version": prompt_version,
"document_id": document_id,
"collection": document_collection,
}
# 5. Guarda respuesta en S3
from datetime import datetime
now = datetime.utcnow()
response_key = f"responses/{now.strftime('%Y/%m/%d')}/{document_id}-{prompt_name}.json"
s3.put_object(
Bucket=BUCKET,
Key=response_key,
Body=json.dumps(result, ensure_ascii=False).encode("utf-8"),
ContentType="application/json",
)
result["response_key"] = response_key
return result
# Uso
output = run_ai_pipeline(
prompt_name="summarizer",
prompt_version="v1",
document_collection="product-docs",
document_id="doc-001",
)
print(f"Respuesta guardada en: s3://{BUCKET}/{output['response_key']}")
Batch Processing con S3 + Lambda
Procesar múltiples documentos
def batch_process_collection(
prompt_name: str,
prompt_version: str,
collection: str,
) -> dict:
"""Procesa todos los documentos de una colección con el prompt indicado."""
paginator = s3.get_paginator("list_objects_v2")
prefix = f"documents/{collection}/"
doc_keys = []
for page in paginator.paginate(Bucket=BUCKET, Prefix=prefix):
for obj in page.get("Contents", []):
if obj["Key"].endswith(".json"):
doc_keys.append(obj["Key"])
print(f"Encontrados {len(doc_keys)} documentos en '{collection}'")
results = {"success": 0, "failed": 0, "errors": []}
for key in doc_keys:
doc_id = key.split("/")[-1].replace(".json", "")
try:
output = run_ai_pipeline(
prompt_name=prompt_name,
prompt_version=prompt_version,
document_collection=collection,
document_id=doc_id,
)
results["success"] += 1
print(f" ✓ {doc_id}: {output['tokens_used']} tokens")
except Exception as e:
results["failed"] += 1
results["errors"].append({"doc_id": doc_id, "error": str(e)})
print(f" ✗ {doc_id}: {str(e)}")
return results
SAM Template con S3 Trigger
template.yaml para S3 → Lambda
AWSTemplateFormatVersion: '2010-09-09'
Transform: AWS::Serverless-2016-10-31
Description: S3 + Lambda AI Pipeline — Module 5
Parameters:
OpenAiApiKey:
Type: String
NoEcho: true
BucketName:
Type: String
Default: ai-assets-dev
Resources:
AIBucket:
Type: AWS::S3::Bucket
Properties:
BucketName: !Ref BucketName
DocProcessorFunction:
Type: AWS::Serverless::Function
Properties:
FunctionName: ai-doc-processor
Handler: s3_trigger_handler.handler
Runtime: python3.11
Architectures: [arm64]
MemorySize: 512
Timeout: 120
Environment:
Variables:
OPENAI_API_KEY: !Ref OpenAiApiKey
AI_BUCKET: !Ref BucketName
PROMPT_KEY: prompts/summarizer/v1/system.txt
RESPONSE_PREFIX: responses/
MODEL_NAME: gpt-4o-mini
Policies:
- S3ReadPolicy:
BucketName: !Ref BucketName
- S3CrudPolicy:
BucketName: !Ref BucketName
Events:
NewDocument:
Type: S3
Properties:
Bucket: !Ref AIBucket
Events: s3:ObjectCreated:*
Filter:
S3Key:
Rules:
- Name: prefix
Value: documents/inbox/
- Name: suffix
Value: .json
Outputs:
BucketName:
Value: !Ref AIBucket
ProcessorArn:
Value: !GetAtt DocProcessorFunction.Arn
Testing del Flujo Completo
Test manual con boto3
import boto3
import json
import os
import time
s3 = boto3.client("s3", endpoint_url=os.environ.get("AWS_ENDPOINT_URL"))
BUCKET = "ai-assets-dev"
def setup_test_data():
"""Prepara datos de test en S3."""
s3.put_object(
Bucket=BUCKET,
Key="prompts/summarizer/v1/system.txt",
Body="Resume el siguiente documento en 3 puntos clave.".encode("utf-8"),
)
test_doc = {
"id": "test-doc-001",
"title": "Introducción a Serverless",
"content": (
"Serverless computing permite ejecutar código sin gestionar servidores. "
"AWS Lambda es el servicio más popular para serverless. "
"Las funciones se ejecutan en respuesta a eventos y escalan automáticamente. "
"Pagas solo por el tiempo de ejecución, no por servidores idle."
),
}
s3.put_object(
Bucket=BUCKET,
Key="documents/inbox/test-doc-001.json",
Body=json.dumps(test_doc, ensure_ascii=False).encode("utf-8"),
)
print("Test data uploaded")
def verify_response():
"""Verifica que Lambda procesó el documento y guardó la respuesta."""
time.sleep(5)
response = s3.list_objects_v2(Bucket=BUCKET, Prefix="responses/")
contents = response.get("Contents", [])
if not contents:
print("No responses found yet. Lambda may still be processing.")
return
for obj in contents:
print(f"Response found: {obj['Key']}")
resp = s3.get_object(Bucket=BUCKET, Key=obj["Key"])
data = json.loads(resp["Body"].read().decode("utf-8"))
print(f" Document: {data.get('document_title')}")
print(f" Answer: {data.get('answer', '')[:200]}...")
print(f" Tokens: {data.get('tokens_used')}")
setup_test_data()
verify_response()
Troubleshooting
Problema 1: Lambda no se triggerear cuando subo archivo a S3
Verifica la configuración de notification y los permisos.
# Verificar notification configuration
aws s3api get-bucket-notification-configuration --bucket ai-assets-dev
# Verificar que Lambda tiene permission policy para S3
aws lambda get-policy --function-name ai-doc-processor
# Si falta el permiso:
aws lambda add-permission \
--function-name ai-doc-processor \
--statement-id s3-trigger \
--action lambda:InvokeFunction \
--principal s3.amazonaws.com \
--source-arn arn:aws:s3:::ai-assets-dev
Problema 2: Loop infinito — Lambda se triggerear a sí misma
Lambda escribe a un prefijo que activa otro trigger.
# Verifica que tu trigger filtra por prefijo de ENTRADA
# y que Lambda escribe a un prefijo de SALIDA diferente
# ✅ Trigger: documents/inbox/*.json
# ✅ Output: responses/2026/03/08/*.json
# ❌ Trigger: documents/*.json (demasiado amplio)
# ❌ Output: documents/processed/*.json (mismo prefijo base)
Problema 3: "NoSuchKey" al leer prompt template
El prompt template no existe en S3 cuando Lambda se ejecuta.
try:
prompt = get_prompt_template()
except s3.exceptions.NoSuchKey:
logger.error(f"Prompt template not found: {PROMPT_KEY}")
prompt = "Resume el siguiente documento de forma concisa."
Problema 4: Evento S3 con key URL-encoded
S3 envía keys URL-encoded. Espacios se convierten en +.
from urllib.parse import unquote_plus
raw_key = record["s3"]["object"]["key"]
key = unquote_plus(raw_key)
Ejercicios Prácticos
Ejercicio 1: Pipeline multi-prompt
Modifica el pipeline para que procese cada documento con múltiples prompts (por ejemplo, "summarizer" y "classifier") y guarde cada resultado en un prefijo separado.
Ver solución
import boto3
import json
import os
from datetime import datetime
from openai import OpenAI
s3 = boto3.client("s3", endpoint_url=os.environ.get("AWS_ENDPOINT_URL"))
openai_client = OpenAI(api_key=os.environ.get("OPENAI_API_KEY", ""))
BUCKET = os.environ.get("AI_BUCKET", "ai-assets-dev")
PIPELINES = [
{"prompt_name": "summarizer", "prompt_version": "v1", "output_prefix": "summaries/"},
{"prompt_name": "classifier", "prompt_version": "v1", "output_prefix": "classifications/"},
]
def multi_prompt_handler(event, context):
results = []
for record in event.get("Records", []):
source_key = record["s3"]["object"]["key"]
source_bucket = record["s3"]["bucket"]["name"]
doc_resp = s3.get_object(Bucket=source_bucket, Key=source_key)
document = json.loads(doc_resp["Body"].read().decode("utf-8"))
for pipeline in PIPELINES:
prompt_key = f"prompts/{pipeline['prompt_name']}/{pipeline['prompt_version']}/system.txt"
prompt_resp = s3.get_object(Bucket=BUCKET, Key=prompt_key)
prompt_template = prompt_resp["Body"].read().decode("utf-8")
full_prompt = f"{prompt_template}\n\n{document.get('content', '')}"
response = openai_client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": full_prompt}],
max_tokens=500,
)
doc_id = document.get("id", source_key.split("/")[-1].replace(".json", ""))
now = datetime.utcnow()
output_key = (
f"{pipeline['output_prefix']}{now.strftime('%Y/%m/%d')}/"
f"{doc_id}-{pipeline['prompt_name']}.json"
)
result = {
"document_id": doc_id,
"pipeline": pipeline["prompt_name"],
"answer": response.choices[0].message.content,
"tokens_used": response.usage.total_tokens,
"processed_at": now.isoformat(),
}
s3.put_object(
Bucket=BUCKET,
Key=output_key,
Body=json.dumps(result, ensure_ascii=False).encode("utf-8"),
ContentType="application/json",
)
results.append({"doc_id": doc_id, "pipeline": pipeline["prompt_name"], "output": output_key})
return {"processed": len(results), "results": results}
Ejercicio 2: S3 trigger con deduplicación
Implementa un mecanismo que evite procesar el mismo documento dos veces. Usa un prefijo "processed/" en S3 como registro de documentos ya procesados.
Ver solución
import boto3
import json
import os
import hashlib
from datetime import datetime
from urllib.parse import unquote_plus
s3 = boto3.client("s3", endpoint_url=os.environ.get("AWS_ENDPOINT_URL"))
BUCKET = os.environ.get("AI_BUCKET", "ai-assets-dev")
PROCESSED_PREFIX = "processed/"
def is_already_processed(source_key: str) -> bool:
"""Verifica si un documento ya fue procesado."""
marker_key = f"{PROCESSED_PREFIX}{hashlib.md5(source_key.encode()).hexdigest()}"
try:
s3.head_object(Bucket=BUCKET, Key=marker_key)
return True
except s3.exceptions.ClientError as e:
if e.response["Error"]["Code"] == "404":
return False
raise
def mark_as_processed(source_key: str, result_key: str):
"""Marca un documento como procesado."""
marker_key = f"{PROCESSED_PREFIX}{hashlib.md5(source_key.encode()).hexdigest()}"
s3.put_object(
Bucket=BUCKET,
Key=marker_key,
Body=json.dumps({
"source_key": source_key,
"result_key": result_key,
"processed_at": datetime.utcnow().isoformat(),
}).encode("utf-8"),
)
def handler(event, context):
results = []
for record in event.get("Records", []):
source_key = unquote_plus(record["s3"]["object"]["key"])
if is_already_processed(source_key):
results.append({"key": source_key, "status": "skipped", "reason": "already processed"})
continue
try:
result_key = process_and_save(source_key)
mark_as_processed(source_key, result_key)
results.append({"key": source_key, "status": "processed", "result": result_key})
except Exception as e:
results.append({"key": source_key, "status": "error", "error": str(e)})
return {"results": results}
def process_and_save(source_key: str) -> str:
"""Procesa documento y retorna key del resultado."""
doc_resp = s3.get_object(Bucket=BUCKET, Key=source_key)
document = json.loads(doc_resp["Body"].read().decode("utf-8"))
now = datetime.utcnow()
doc_id = source_key.split("/")[-1].replace(".json", "")
result_key = f"responses/{now.strftime('%Y/%m/%d')}/{doc_id}.json"
s3.put_object(
Bucket=BUCKET,
Key=result_key,
Body=json.dumps({"processed": True, "source": source_key}).encode("utf-8"),
)
return result_key
Ejercicio 3: Pipeline con versionado de prompts A/B
Implementa un flujo que procese cada documento con dos versiones de prompt (A/B test) y guarde ambos resultados para comparar.
Ver solución
import boto3
import json
import os
import time
from datetime import datetime
from openai import OpenAI
s3 = boto3.client("s3", endpoint_url=os.environ.get("AWS_ENDPOINT_URL"))
openai_client = OpenAI(api_key=os.environ.get("OPENAI_API_KEY", ""))
BUCKET = os.environ.get("AI_BUCKET", "ai-assets-dev")
def ab_test_handler(event, context):
"""Procesa cada documento con dos versiones de prompt para A/B testing."""
versions = ["v1", "v2"]
prompt_name = os.environ.get("PROMPT_NAME", "summarizer")
results = []
for record in event.get("Records", []):
source_key = record["s3"]["object"]["key"]
doc_resp = s3.get_object(Bucket=BUCKET, Key=source_key)
document = json.loads(doc_resp["Body"].read().decode("utf-8"))
doc_id = document.get("id", "unknown")
ab_results = {}
for version in versions:
prompt_key = f"prompts/{prompt_name}/{version}/system.txt"
try:
prompt_resp = s3.get_object(Bucket=BUCKET, Key=prompt_key)
prompt_template = prompt_resp["Body"].read().decode("utf-8")
except Exception:
continue
full_prompt = f"{prompt_template}\n\n{document.get('content', '')}"
start = time.time()
response = openai_client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": full_prompt}],
max_tokens=500,
)
duration = int((time.time() - start) * 1000)
ab_results[version] = {
"answer": response.choices[0].message.content,
"tokens_used": response.usage.total_tokens,
"duration_ms": duration,
}
now = datetime.utcnow()
comparison_key = f"ab-tests/{prompt_name}/{now.strftime('%Y/%m/%d')}/{doc_id}.json"
s3.put_object(
Bucket=BUCKET,
Key=comparison_key,
Body=json.dumps({
"document_id": doc_id,
"prompt_name": prompt_name,
"versions_tested": versions,
"results": ab_results,
"tested_at": now.isoformat(),
}, ensure_ascii=False).encode("utf-8"),
ContentType="application/json",
)
results.append({"doc_id": doc_id, "comparison_key": comparison_key})
return {"ab_tests": len(results), "results": results}
Ejercicio 4: Monitor de pipeline con alertas
Crea una función que revise el bucket de responses y genere un reporte de salud del pipeline: documentos procesados en las últimas 24h, errores, tokens promedio, etc.
Ver solución
import boto3
import json
import os
from datetime import datetime, timedelta
s3 = boto3.client("s3", endpoint_url=os.environ.get("AWS_ENDPOINT_URL"))
BUCKET = os.environ.get("AI_BUCKET", "ai-assets-dev")
def pipeline_health_check() -> dict:
"""Genera reporte de salud del pipeline de las últimas 24h."""
now = datetime.utcnow()
yesterday = now - timedelta(days=1)
prefixes = [
f"responses/{now.strftime('%Y/%m/%d')}/",
f"responses/{yesterday.strftime('%Y/%m/%d')}/",
]
all_responses = []
for prefix in prefixes:
paginator = s3.get_paginator("list_objects_v2")
for page in paginator.paginate(Bucket=BUCKET, Prefix=prefix):
for obj in page.get("Contents", []):
if obj["LastModified"].replace(tzinfo=None) >= yesterday:
try:
resp = s3.get_object(Bucket=BUCKET, Key=obj["Key"])
data = json.loads(resp["Body"].read().decode("utf-8"))
all_responses.append(data)
except Exception:
pass
if not all_responses:
return {"status": "warning", "message": "No responses in last 24h", "count": 0}
tokens_list = [r.get("tokens_used", 0) for r in all_responses]
errors = [r for r in all_responses if r.get("status") == "error"]
report = {
"status": "healthy" if len(errors) == 0 else "degraded",
"period": f"{yesterday.isoformat()} — {now.isoformat()}",
"total_processed": len(all_responses),
"errors": len(errors),
"error_rate": f"{(len(errors) / len(all_responses)) * 100:.1f}%",
"tokens": {
"total": sum(tokens_list),
"average": round(sum(tokens_list) / len(tokens_list)),
"max": max(tokens_list),
},
"checked_at": now.isoformat(),
}
if len(errors) / max(len(all_responses), 1) > 0.1:
report["alert"] = "Error rate > 10%. Check pipeline configuration."
return report
health = pipeline_health_check()
print(json.dumps(health, indent=2))
Resumen
- S3 Event Notifications triggerean Lambda automáticamente cuando se sube un archivo. Configura filtros de prefijo y sufijo para precisión.
- El flujo S3 → Lambda → S3 es el patrón central de un servicio AI en AWS: Lambda lee assets (prompts, documentos), invoca el LLM, y persiste resultados.
- Evita loops infinitos usando prefijos de entrada y salida separados, o buckets diferentes para input y output.
- Parsea el evento S3 correctamente: decodifica URL-encoded keys con
unquote_plus, maneja múltiples Records. - SAM Template define el trigger S3 → Lambda como código — reproducible y versionable.
- Batch processing permite procesar colecciones completas de documentos con un solo prompt o múltiples prompts (A/B testing).
Recursos Adicionales
- S3 Event Notifications — Configuración de triggers
- Lambda S3 Event Tutorial — Tutorial oficial
- S3 Event Message Structure — Formato del evento
- SAM S3 Event — SAM template para S3 triggers
- Lambda Permissions for S3 — Permisos requeridos
- Avoiding Recursive Invocations — Prevenir loops
- S3 Batch Operations — Procesamiento masivo nativo
- Lambda Dead Letter Queues — Manejo de fallos asíncronos