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

  1. S3 Event Notifications — Configuración de triggers
  2. Lambda S3 Event Tutorial — Tutorial oficial
  3. S3 Event Message Structure — Formato del evento
  4. SAM S3 Event — SAM template para S3 triggers
  5. Lambda Permissions for S3 — Permisos requeridos
  6. Avoiding Recursive Invocations — Prevenir loops
  7. S3 Batch Operations — Procesamiento masivo nativo
  8. Lambda Dead Letter Queues — Manejo de fallos asíncronos