La lección anterior terminó con dos preguntas. La primera: ¿hace falta un Deployment en EKS, con sus réplicas, sus sondas, su HPA y su guardia, para generar tres miniaturas cada vez que la Quesería Montblanc sube una foto, o para producir un PDF de factura cuando llega un pago.confirmado? Son tareas pequeñas, esporádicas, sin estado, que responden a un evento y terminan. La segunda: ¿por qué la foto de queso-curado tiene que viajar desde la región eu-west-1 hasta el móvil de Ana en Valencia en cada visita, pagando egress y 40 ms de latencia, cuando podría estar guardada a veinte kilómetros de ella? Y hay una tercera, que ya asomó en 08-02: ¿qué hace furgoneta-3 con sus posiciones durante un túnel de tres minutos, o el puesto de la Huerta La Vega en el mercado de Lleida cuando la conexión se cae y siguen llegando clientes? Las tres preguntas tienen respuestas que salen del modelo "servicio en un clúster en una región": serverless, donde el proveedor ejecuta funciones efímeras en respuesta a eventos y cobra por invocación; y edge computing, donde el cómputo y los datos se acercan a quien los usa, sea una CDN, una función en el borde de la red o un dispositivo que funciona sin conexión. Esta lección desarrolla ambos, con su modelo de ejecución, sus límites, sus casos adecuados e inadecuados en Kilómetro Cero, y el código que los implementa: una Lambda de miniaturas, una de facturas, la saga de pedido reescrita en Step Functions, un Worker de Cloudflare que verifica el JWT en el borde y un agregador local para la furgoneta. El proyecto final (08-05) integrará todo.
Contenido
- Serverless y FaaS: el modelo de ejecución
- Arranque en frío, límites, concurrencia y precio
- Casos adecuados e inadecuados: los cuatro de Kilómetro Cero
- Serverless más allá de FaaS: colas, bases de datos, API Gateway y Step Functions
- Patrones: fan-out, idempotencia obligatoria y DLQ
- Observabilidad, pruebas locales, lock-in y coste a gran escala
- Código serverless: miniaturas, facturas, SAM y la saga en Step Functions
- Edge computing: por qué acercar el cómputo
- CDN: cachear fotos y catálogo cerca de Ana
- Funciones en el borde: verificación de JWT y personalización
- Edge para IoT: la furgoneta y el mercado sin conexión
- El continuo nube-edge-dispositivo y la seguridad en el borde
- Errores Comunes y Consejos
- Ejercicios
- Conclusión
- Serverless y FaaS: el modelo de ejecución
"Serverless" no significa que no haya servidores: significa que no son nuestros ni los vemos. El proveedor los aprovisiona, los escala y los retira; nosotros entregamos código y pagamos por lo que se ejecuta. La forma más pura es FaaS (Function as a Service): AWS Lambda, Google Cloud Functions, Azure Functions. En 08-03 la tabla de responsabilidad compartida lo dejaba nombrado; este es su modelo de ejecución:
- Un evento ocurre: un objeto llega a un bucket, un mensaje entra en una cola o un tópico, una petición HTTP llega a un endpoint, un temporizador salta.
- El proveedor busca una instancia de la función lista (un microcontenedor con el runtime de Python y nuestro código). Si no hay ninguna libre, crea una (arranque en frío).
- Invoca el manejador (
handler(evento, contexto)) con el evento serializado. La función hace su trabajo y devuelve. - La instancia queda congelada unos minutos por si llega otro evento; si no llega, se destruye. No hay estado entre invocaciones en el que se pueda confiar (el sistema de ficheros temporal puede sobrevivir, pero no está garantizado).
- Si llegan mil eventos a la vez, el proveedor crea hasta mil instancias en paralelo (escalado a miles); si no llega ninguno, no hay ninguna (escalado a cero) y no se paga nada.
La comparación con el Deployment de 07-05 resume la diferencia de modelo:
| Aspecto | Servicio en Kubernetes (07-05) | Función serverless |
|---|---|---|
| Unidad | Contenedor de larga vida con varias réplicas | Función efímera, una invocación por instancia (en Lambda) |
| Escalado | HPA por métrica, de n a m réplicas, en decenas de segundos | Automático por evento, de 0 a miles, en milisegundos o segundos |
| Coste en reposo | Las réplicas mínimas se pagan siempre | Cero |
| Estado | En memoria mientras vive el Pod; en Redis o base de datos | Solo externo (S3, DynamoDB, base de datos) |
| Conexiones | Pool de conexiones persistente a PostgreSQL, Kafka | Cada instancia abre las suyas; miles de instancias agotan las conexiones de una base de datos (hace falta un proxy: RDS Proxy) |
| Duración | Ilimitada (un consumidor de Kafka corre días) | Limitada (15 min en Lambda) |
| Despliegue | Imagen + manifiesto + rollout | Paquete de código (zip o imagen) + definición del evento |
| Operación | Sondas, recursos, guardia por SLO | Sin servidores que operar; sí métricas, errores y DLQ que vigilar |
- Arranque en frío, límites, concurrencia y precio
Arranque en frío (cold start): crear la instancia, cargar el runtime y ejecutar el código de inicialización (imports, conexiones) antes de la primera invocación. En Python con pocas dependencias son 100-300 ms; con Pillow, boto3 y un cliente de Kafka, 1-2 s; en una VPC, algo más. Es aceptable para procesar una foto y no para responder a Ana. Se mitiga con concurrencia aprovisionada (instancias mantenidas calientes, que se pagan aunque no se usen: se pierde el escalado a cero), con paquetes pequeños, y moviendo la inicialización pesada fuera del manejador para que se ejecute una vez por instancia y no por invocación.
Límites: tiempo máximo por invocación (15 min en Lambda, 9-60 min en Cloud Functions según la generación), memoria (de 128 MB a 10 GB; la CPU se asigna en proporción a la memoria, así que una función lenta a veces se acelera dándole memoria), tamaño del paquete (250 MB descomprimido; o una imagen de contenedor de hasta 10 GB), tamaño del evento (6 MB síncrono, 256 KB asíncrono), espacio temporal (/tmp de hasta 10 GB).
Concurrencia: el número de instancias simultáneas. Tiene un límite por cuenta y región (1 000 por defecto en Lambda, ampliable) y se puede reservar por función (para que una función descontrolada no agote la cuota de las demás) o limitar (para no agotar la base de datos de detrás: si km0_analitica acepta 100 conexiones, la función que la escribe no puede tener 500 instancias). Ese límite es el bulkhead de 07-04 en versión serverless.
Precio: por número de invocaciones (del orden de 0,20 € por millón) más por GB-segundo de ejecución (memoria asignada × duración; del orden de 0,000017 € por GB-s), con un tramo gratuito mensual. Una función de miniaturas con 512 MB que tarda 800 ms cuesta ≈ 0,000007 € por foto: 100 000 fotos al mes, menos de 1 €. Es lo que hace que serverless sea imbatible para cargas esporádicas y lo que lo hace caro para cargas constantes (apartado 6).
- Casos adecuados e inadecuados: los cuatro de Kilómetro Cero
| Característica de la carga | Adecuado para FaaS | Inadecuado para FaaS |
|---|---|---|
| Frecuencia | Esporádica o con picos enormes e imprevisibles | Constante y alta (un consumidor que procesa 500 eventos/s todo el día) |
| Duración | Segundos, como mucho minutos | Horas (un trabajo de Spark; un consumidor de Kafka de larga vida) |
| Estado | Sin estado; todo en servicios externos | Estado en memoria entre peticiones (la tabla de suscripciones WebSocket de 08-02) |
| Latencia | Tolera cientos de ms (procesamiento en segundo plano) | p99 < 500 ms exigido a personas (el SLO de pedidos) |
| Conexiones | Pocas, breves, a servicios que escalan (S3, DynamoDB, SQS) | Pool persistente a PostgreSQL, sesiones WebSocket, conexiones MQTT |
| Disparo | Un evento con nombre (objeto, mensaje, petición, temporizador) | Un bucle que sondea o un proceso que escucha un puerto |
Con ese filtro, Kilómetro Cero identifica cuatro tareas que hoy viven como código dentro de los servicios (o como cron jobs en Kubernetes) y que encajan mejor como funciones:
- Miniaturas al subir una foto a
km0-fotos. En 04-03 la Quesería Montblanc subíafotos/queso-curado/original.jpgpor URL prefirmada y "un proceso de miniaturas" generaba las de 300 y 800 px. Ese proceso era un consumidor con un bucle. Como función: eventos3:ObjectCreated→ Lambda con Pillow → escribeminiatura-300.jpgyminiatura-800.jpg. Esporádico (cientos de fotos al día, con picos cuando un productor sube el catálogo entero), sin estado, 1 s por foto. - Factura PDF al
pago.confirmado. Hoy un consumidor depedidos.eventosenpedidos. Como función: evento de Kafka (MSK como origen de eventos de Lambda, en lotes) → generar el PDF → escribirlo enkm0-facturas→ notificar. Ritmo de los pedidos: miles al día, picos en campaña, sin estado. - Webhooks de la pasarela de pagos. La pasarela notifica de forma asíncrona los cobros diferidos y las devoluciones con un
POSTa una URL nuestra. Es un endpoint que recibe pocas llamadas, debe estar siempre disponible, verificar una firma y publicar un evento: API Gateway + Lambda, sin ocupar una réplica depagospara esperar. - Tareas programadas ligeras. Caducar reservas de stock no confirmadas cada minuto, limpiar sesiones, comprobar que el DAG de Airflow terminó: temporizador (EventBridge) → Lambda. Sustituyen a CronJobs de Kubernetes que ocupaban recursos para ejecutar 200 ms de trabajo.
Y lo que no se mueve a funciones, con la razón: pedidos y catalogo (latencia exigida a personas, pool de conexiones, Kong y mesh delante), el servidor WebSocket de 08-02 (estado y conexiones largas; hay servicios gestionados de WebSocket, pero el fan-out entre instancias y la lógica de suscripción son justo lo que no encaja), los consumidores de Flink (estado y checkpoints), los trabajos de Spark (horas).
- Serverless más allá de FaaS: colas, bases de datos, API Gateway y Step Functions
FaaS es la parte visible; a su alrededor hay servicios "serverless" en el sentido de sin capacidad que aprovisionar y con pago por uso:
| Servicio | Qué es | Equivalente que Kilómetro Cero ya tiene | Cuándo preferirlo |
|---|---|---|---|
| Colas (SQS, Cloud Tasks) | Cola gestionada, escalado ilimitado, pago por mensaje, con DLQ integrada | RabbitMQ (02-04) | Para conectar funciones entre sí y absorber picos sin operar un broker |
| Bases de datos serverless (DynamoDB, Aurora Serverless, Firestore) | Clave-valor o relacional que escala capacidad automáticamente, pago por petición o por capacidad consumida | Cassandra (04-04), RDS (08-03) | DynamoDB para el estado de las funciones (tabla de idempotencia, checkpoints); Aurora Serverless para bases con carga muy variable |
| API Gateway gestionado | Endpoint HTTP que enruta a funciones, con autenticación (JWT), límites y claves de API | Kong (06-05) | Para exponer funciones (webhooks) sin pasar por el clúster; con menos plugins que Kong |
| Step Functions / Workflows | Orquestador de flujos con estado: define pasos, ramas, reintentos y compensaciones en JSON, y ejecuta cada paso invocando funciones o servicios | saga_pedido.py con la tabla sagas (03-05) |
Cuando la saga se compone de funciones y no se quiere escribir ni operar el orquestador |
Step Functions merece atención porque sustituye a algo que el curso construyó con esfuerzo. El OrquestadorSagaPedido de 03-05 persistía el estado de cada saga en la tabla sagas, reintentaba los fallos transitorios, distinguía los permanentes y ejecutaba las compensaciones en orden inverso. Step Functions hace exactamente eso como servicio: el estado de cada ejecución lo guarda el proveedor (con historial completo de cada paso), los reintentos y los Catch se declaran por paso, y las compensaciones son ramas del grafo. Lo que se pierde: el orquestador vive fuera del clúster, cada paso es una invocación (con su latencia y su coste), la definición es JSON y no Python, y es lo más específico del proveedor que hay (lock-in total). La comparación completa está en el apartado 7, con la saga reescrita.
- Patrones: fan-out, idempotencia obligatoria y DLQ
Fan-out con colas. Un evento que debe procesarse de varias formas (una foto subida: miniaturas, análisis de contenido, actualización de la ficha) no dispara tres funciones directamente desde S3; se publica en un tópico (SNS o EventBridge) que lo entrega a una cola por consumidor (SQS), y cada cola dispara su función. Las colas desacoplan (si la función de análisis falla, las miniaturas se generan igual), absorben picos (un productor que sube 2 000 fotos no crea 2 000 × 3 instancias de golpe: la cola las dosifica con el límite de concurrencia) y dan reintentos y DLQ.
Idempotencia obligatoria. El proveedor entrega los eventos al menos una vez y reintenta automáticamente cuando la función falla o agota el tiempo: una función que tarda 16 minutos por un pico de tamaño se reinvoca, y si la primera instancia había escrito la mitad del resultado, la segunda parte de ese estado. Todo lo que se vio en 02-05 sobre consumidores idempotentes aplica con más fuerza, porque aquí los reintentos no los controlamos nosotros: cada función debe ser segura ante repetición. Las técnicas: claves de salida deterministas (la miniatura se escribe siempre en la misma clave: escribir dos veces es inofensivo), tabla de idempotencia en DynamoDB con el id del evento (la factura F-2026-000124 se genera una vez), y operaciones condicionales (PutItem con attribute_not_exists).
DLQ. Tras n reintentos (configurable por función o por cola), el evento va a una cola de mensajes muertos, con alerta (07-01) y procedimiento de reproceso: la DLQ de 02-05, gestionada. Sin DLQ configurada, un evento envenenado (una imagen corrupta que hace fallar a Pillow) se reintenta hasta agotar la política y se pierde en silencio.
- Observabilidad, pruebas locales, lock-in y coste a gran escala
Observabilidad. Las funciones emiten logs a CloudWatch (o Cloud Logging) automáticamente, y métricas de invocaciones, errores, duración y throttles. Lo que hay que añadir es lo mismo que en 07-02: logs estructurados con el X-Request-Id o el id_evento para correlacionar, y trazas con OpenTelemetry (o X-Ray) que enlacen la función con el servicio que produjo el evento, porque una petición de Ana que acaba en una factura atraviesa pedidos, Kafka y la Lambda, y sin propagación de contexto la traza se corta en Kafka. El servicios/comun/trazas.py de 07-02 se empaqueta con la función como capa (layer).
Pruebas locales. Una función no se puede ejecutar "sin el proveedor" salvo con emuladores: AWS SAM (sam local invoke, sam local start-api) ejecuta la función en un contenedor con un evento de ejemplo; LocalStack emula S3, SQS, DynamoDB y más en Docker, lo que permite pruebas de integración con Testcontainers (07-06) sin cuenta de AWS. Ninguno es perfecto (los permisos IAM y los límites reales solo se ven en el proveedor), así que el pipeline despliega a un entorno de staging real para las pruebas de contrato.
Lock-in. Es el más alto de todo lo visto: el formato del evento, el manejador, los servicios de alrededor (SQS, DynamoDB, Step Functions) son del proveedor. Se mitiga separando el manejador (adaptador de 10 líneas) de la lógica (funciones Python puras, probables sin proveedor), como se hace en el apartado 7; pero hay que aceptarlo: quien elige Step Functions elige AWS.
Coste a gran escala. El precio por invocación es imbatible a baja frecuencia y se vuelve caro a alta: una función con 512 MB que corre constantemente (por ejemplo, 100 invocaciones/s de 200 ms cada una, 24 h al día) consume ≈ 2,6 millones de GB-s al mes ≈ 45 €, más 260 millones de invocaciones ≈ 52 €: unos 100 €/mes, frente a ≈ 30 € de un Pod pequeño en un nodo ya pagado. La regla aproximada: cuando la utilización media supera el 30-40 % de un contenedor equivalente, el contenedor sale más barato. Por eso las miniaturas y las facturas son serverless y pedidos no.
- Código serverless: miniaturas, facturas, SAM y la saga en Step Functions
La estructura añadida a km0/:
km0/serverless/
├── template.yaml # AWS SAM: funciones, eventos, permisos, DLQ
├── miniaturas/
│ ├── handler.py
│ └── requirements.txt # Pillow, boto3
├── facturas/
│ ├── handler.py
│ └── requirements.txt # reportlab, boto3
└── saga/
└── saga_pedido.asl.json # Step Functions (Amazon States Language)El flujo S3 → Lambda
flowchart LR
Q[Quesería Montblanc<br/>PUT por URL prefirmada] --> S3[(S3 km0-fotos<br/>fotos/queso-curado/original.jpg)]
S3 -- "s3:ObjectCreated:*<br/>prefijo fotos/, sufijo original.jpg" --> L[Lambda miniaturas<br/>Pillow, 1024 MB, 60 s]
L -- "PUT miniatura-300.jpg<br/>PUT miniatura-800.jpg" --> S3
L -- "evento foto.procesada" --> EB[EventBridge]
EB --> CAT[catalogo: actualiza ficha]
L -. "fallo tras 2 reintentos" .-> DLQ[(SQS DLQ miniaturas)]
DLQ -. alerta .-> OPS[Alerta 07-01]
serverless/miniaturas/handler.py
# km0/serverless/miniaturas/handler.py
"""Genera miniaturas de 300 y 800 px al subir un original a km0-fotos.
Idempotente: las claves de salida son deterministas y el ETag del original se guarda
como metadato de cada miniatura; si ya coincide, no se regenera.
"""
import io, json, os, urllib.parse
import boto3
from PIL import Image
# Inicialización FUERA del manejador: se ejecuta una vez por instancia (arranque en frío),
# no una vez por invocación. Aquí van los clientes y todo lo costoso.
s3 = boto3.client("s3")
eventbridge = boto3.client("events")
TAMANOS = (300, 800)
BUS_EVENTOS = os.environ.get("KM0_BUS_EVENTOS", "km0")
def _clave_miniatura(clave_original: str, px: int) -> str:
# fotos/queso-curado/original.jpg -> fotos/queso-curado/miniatura-300.jpg (convención de 04-03)
prefijo = clave_original.rsplit("/", 1)[0]
return f"{prefijo}/miniatura-{px}.jpg"
def _ya_generada(bucket: str, clave: str, etag_original: str) -> bool:
try:
cabecera = s3.head_object(Bucket=bucket, Key=clave)
return cabecera.get("Metadata", {}).get("origen-etag") == etag_original
except s3.exceptions.ClientError:
return False
def procesar_objeto(bucket: str, clave: str, etag_original: str) -> list[str]:
"""Lógica pura de negocio: probable sin AWS pasando un cliente falso en s3."""
cuerpo = s3.get_object(Bucket=bucket, Key=clave)["Body"].read()
imagen = Image.open(io.BytesIO(cuerpo))
imagen = imagen.convert("RGB") # PNG con alfa -> JPEG
generadas = []
for px in TAMANOS:
destino = _clave_miniatura(clave, px)
if _ya_generada(bucket, destino, etag_original): # reintento o duplicado: no rehacer
generadas.append(destino)
continue
copia = imagen.copy()
copia.thumbnail((px, px)) # conserva la proporción; nunca amplía
salida = io.BytesIO()
copia.save(salida, format="JPEG", quality=85, optimize=True)
s3.put_object(
Bucket=bucket, Key=destino, Body=salida.getvalue(), ContentType="image/jpeg",
CacheControl="public, max-age=31536000, immutable", # la CDN (apartado 9) la cachea un año
Metadata={"origen-etag": etag_original}, # marca de idempotencia
)
generadas.append(destino)
return generadas
def handler(evento: dict, contexto) -> dict:
"""Adaptador: traduce el evento de S3 a llamadas a la lógica. 10 líneas, sin negocio."""
resultados = []
for registro in evento["Records"]: # S3 puede agrupar varios objetos en un evento
bucket = registro["s3"]["bucket"]["name"]
clave = urllib.parse.unquote_plus(registro["s3"]["object"]["key"]) # las claves llegan URL-codificadas
etag = registro["s3"]["object"].get("eTag", "")
if not clave.endswith("/original.jpg"):
continue # defensa: el filtro de SAM ya lo hace, pero no fiarse
generadas = procesar_objeto(bucket, clave, etag)
eventbridge.put_events(Entries=[{
"Source": "km0.fotos", "DetailType": "foto.procesada", "EventBusName": BUS_EVENTOS,
"Detail": json.dumps({"bucket": bucket, "original": clave, "miniaturas": generadas}),
}])
print(json.dumps({"nivel": "info", "msg": "miniaturas generadas", "clave": clave,
"n": len(generadas), "request_id": contexto.aws_request_id})) # log estructurado (07-02)
resultados.append(clave)
return {"procesadas": resultados}Decisiones a destacar: la idempotencia se basa en el ETag del original guardado como metadato de la miniatura, de modo que una reinvocación no rehace trabajo y una foto nueva con la misma clave (nuevo ETag) sí lo rehace; CacheControl: immutable es la promesa de 04-03 de no sobrescribir claves cacheadas por la CDN, que aquí se cumple porque la web referencia miniatura-800.jpg?v=<etag>; y el manejador no contiene negocio, para que procesar_objeto se pruebe con un cliente falso y para que el lock-in se limite a esas líneas.
serverless/facturas/handler.py
# km0/serverless/facturas/handler.py
"""Genera la factura PDF de un pedido al recibir pago.confirmado desde MSK (Kafka).
Lambda recibe LOTES de registros por partición. Idempotente por id_evento con DynamoDB
(tabla km0-idempotencia): la misma factura nunca se genera ni envía dos veces.
"""
import base64, io, json, os, time
import boto3
from botocore.exceptions import ClientError
from reportlab.lib.pagesizes import A4
from reportlab.pdfgen import canvas
s3 = boto3.client("s3")
dynamo = boto3.resource("dynamodb").Table(os.environ["TABLA_IDEMPOTENCIA"])
BUCKET_FACTURAS = os.environ["BUCKET_FACTURAS"] # km0-facturas
TTL_S = 30 * 24 * 3600 # la marca de idempotencia caduca a los 30 días
def reclamar(id_evento: str) -> bool:
"""Intenta registrar el evento. Devuelve False si ya estaba (duplicado). Atómico en DynamoDB."""
try:
dynamo.put_item(Item={"id": f"factura#{id_evento}", "ttl": int(time.time()) + TTL_S},
ConditionExpression="attribute_not_exists(id)")
return True
except ClientError as e:
if e.response["Error"]["Code"] == "ConditionalCheckFailedException":
return False
raise
def generar_pdf(pedido: dict) -> bytes:
"""Lógica pura: un PDF sencillo con reportlab."""
buf = io.BytesIO()
c = canvas.Canvas(buf, pagesize=A4)
c.setFont("Helvetica-Bold", 16); c.drawString(50, 800, "Kilómetro Cero - Factura")
c.setFont("Helvetica", 11)
c.drawString(50, 775, f"Factura: F-{pedido['pedido_id'][2:]} Pedido: {pedido['pedido_id']}")
c.drawString(50, 760, f"Cliente: {pedido['cliente']} Mercado: {pedido['mercado']}")
y = 730
for linea in pedido["lineas"]:
c.drawString(60, y, f"{linea['unidades']} x {linea['producto']} ({linea['productor']})")
c.drawRightString(540, y, f"{linea['importe_cents'] / 100:.2f} EUR"); y -= 16
c.setFont("Helvetica-Bold", 12); c.drawRightString(540, y - 10, f"Total: {pedido['total_cents'] / 100:.2f} EUR")
c.showPage(); c.save()
return buf.getvalue()
def handler(evento: dict, contexto) -> dict:
generadas, duplicados = 0, 0
# Formato del evento de MSK: {"records": {"pedidos.eventos-3": [ {value: <base64>, ...}, ... ]}}
for particion, registros in evento["records"].items():
for r in registros:
ev = json.loads(base64.b64decode(r["value"]))
if ev["tipo"] != "pago.confirmado":
continue
if not reclamar(ev["id_evento"]): # reintento de Lambda o duplicado de Kafka (02-05)
duplicados += 1
continue
pedido = ev["datos"]
clave = f"facturas/{pedido['cliente']}/{pedido['pedido_id']}.pdf" # clave determinista
s3.put_object(Bucket=BUCKET_FACTURAS, Key=clave, Body=generar_pdf(pedido),
ContentType="application/pdf", ServerSideEncryption="aws:kms") # cifrado (06-02)
generadas += 1
print(json.dumps({"nivel": "info", "msg": "lote procesado", "generadas": generadas,
"duplicados": duplicados, "request_id": contexto.aws_request_id}))
return {"generadas": generadas}Un matiz importante sobre el orden de operaciones: reclamar se hace antes de generar el PDF. Si la función muere después de reclamar y antes de escribir, la factura no se generará nunca en un reintento (la marca ya existe). La alternativa (reclamar después de escribir) permite duplicados si muere entre medias. Como escribir dos veces la misma clave de S3 es inofensivo (idempotente por sí mismo), la elección correcta aquí es reclamar después, o usar la existencia del objeto como marca. Se deja como está a propósito para el ejercicio 1, que pide corregirlo: es el tipo de razonamiento sobre el punto de ack que 02-05 enseñó, y en serverless no hay excusa para no hacerlo.
serverless/template.yaml con AWS SAM
# km0/serverless/template.yaml
AWSTemplateFormatVersion: "2010-09-09"
Transform: AWS::Serverless-2016-10-31
Description: Funciones de evento de Kilómetro Cero (miniaturas, facturas)
Globals:
Function:
Runtime: python3.12
Architectures: [arm64] # Graviton: más barato por GB-s
Tracing: Active # trazas X-Ray / OpenTelemetry (07-02)
Environment:
Variables: { KM0_ENTORNO: prod }
Resources:
DlqMiniaturas:
Type: AWS::SQS::Queue
Properties: { QueueName: km0-miniaturas-dlq, MessageRetentionPeriod: 1209600 } # 14 días para reprocesar
Miniaturas:
Type: AWS::Serverless::Function
Properties:
CodeUri: miniaturas/
Handler: handler.handler
MemorySize: 1024 # Pillow es CPU: más memoria = más CPU = menos duración
Timeout: 60
ReservedConcurrentExecutions: 50 # bulkhead: un productor que sube 2 000 fotos no agota la cuenta
DeadLetterQueue: { Type: SQS, TargetArn: !GetAtt DlqMiniaturas.Arn }
EventInvokeConfig: { MaximumRetryAttempts: 2 } # invocación asíncrona: 2 reintentos y a la DLQ
Policies:
- S3CrudPolicy: { BucketName: km0-fotos-prod } # solo este bucket
- EventBridgePutEventsPolicy: { EventBusName: km0 }
Events:
FotoSubida:
Type: S3
Properties:
Bucket: !Ref BucketFotos
Events: s3:ObjectCreated:*
Filter:
S3Key:
Rules:
- { Name: prefix, Value: fotos/ }
- { Name: suffix, Value: original.jpg } # NO se dispara con las miniaturas: evita el bucle infinito
BucketFotos: # referenciado desde Terraform (08-03) o creado aquí en staging
Type: AWS::S3::Bucket
Properties: { BucketName: km0-fotos-prod }
TablaIdempotencia:
Type: AWS::DynamoDB::Table
Properties:
TableName: km0-idempotencia
BillingMode: PAY_PER_REQUEST # serverless: sin capacidad que aprovisionar
AttributeDefinitions: [{ AttributeName: id, AttributeType: S }]
KeySchema: [{ AttributeName: id, KeyType: HASH }]
TimeToLiveSpecification: { AttributeName: ttl, Enabled: true }
Facturas:
Type: AWS::Serverless::Function
Properties:
CodeUri: facturas/
Handler: handler.handler
MemorySize: 512
Timeout: 120
Environment:
Variables: { BUCKET_FACTURAS: km0-facturas-prod, TABLA_IDEMPOTENCIA: !Ref TablaIdempotencia }
Policies:
- S3WritePolicy: { BucketName: km0-facturas-prod }
- DynamoDBCrudPolicy: { TableName: !Ref TablaIdempotencia }
- KMSEncryptPolicy: { KeyId: !ImportValue km0-kms-datos }
VpcConfig: # MSK está en las subredes de datos de la VPC (08-03)
SubnetIds: !Split [",", !ImportValue km0-subredes-datos]
SecurityGroupIds: [!ImportValue km0-sg-lambda-msk]
Events:
PagoConfirmado:
Type: MSK
Properties:
Stream: !ImportValue km0-msk-arn
Topics: [pedidos.eventos]
StartingPosition: LATEST
BatchSize: 50 # hasta 50 registros por invocación
MaximumBatchingWindowInSeconds: 5
ConsumerGroupId: facturas-lambda # un grupo de consumidores más en el tópico (02-04)
DestinationConfig:
OnFailure: { Destination: !GetAtt DlqFacturas.Arn }
DlqFacturas:
Type: AWS::SQS::Queue
Properties: { QueueName: km0-facturas-dlq, MessageRetentionPeriod: 1209600 }Tres detalles del template.yaml que evitan errores clásicos: el filtro por sufijo original.jpg impide que la escritura de una miniatura dispare de nuevo la función (un bucle infinito que también sería una factura infinita); ReservedConcurrentExecutions acota el daño de un pico; y el origen de eventos MSK convierte a Lambda en un grupo de consumidores más de pedidos.eventos, con el mismo modelo de particiones y offsets de 02-04, lo que significa que un fallo persistente de la función bloquea el avance de esa partición para ese grupo (y solo para ese grupo) hasta que el lote va a la DLQ. Se despliega con sam build && sam deploy --guided la primera vez, y desde el pipeline después.
La saga de pedido en Step Functions
La saga de 03-05 (reservar stock → cobrar → confirmar; compensar en orden inverso ante fallo permanente), expresada en Amazon States Language, con cada paso invocando una función (o, en producción, el servicio correspondiente vía API):
{
"Comment": "Saga de pedido de Kilómetro Cero (equivalente a saga_pedido.py, 03-05)",
"StartAt": "ReservarStock",
"States": {
"ReservarStock": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": {
"FunctionName": "km0-inventario-reservar",
"Payload": { "id_reserva.$": "$.pedido_id", "pedido_id.$": "$.pedido_id",
"mercado.$": "$.mercado", "lineas.$": "$.lineas" }
},
"ResultPath": "$.reserva",
"Retry": [
{ "ErrorEquals": ["FalloTransitorio", "Lambda.ServiceException", "Lambda.TooManyRequestsException"],
"IntervalSeconds": 1, "MaxAttempts": 3, "BackoffRate": 2.0 }
],
"Catch": [
{ "ErrorEquals": ["StockInsuficiente"], "ResultPath": "$.error", "Next": "RechazarPedido" },
{ "ErrorEquals": ["States.ALL"], "ResultPath": "$.error", "Next": "RechazarPedido" }
],
"Next": "Cobrar"
},
"Cobrar": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": {
"FunctionName": "km0-pagos-cobrar",
"Payload": { "idempotency_key.$": "States.Format('cobro-{}', $.pedido_id)",
"cliente.$": "$.cliente", "importe_cents.$": "$.total_cents" }
},
"ResultPath": "$.cobro",
"TimeoutSeconds": 10,
"Retry": [
{ "ErrorEquals": ["FalloTransitorio", "States.Timeout"], "IntervalSeconds": 2, "MaxAttempts": 3, "BackoffRate": 2.0 }
],
"Catch": [
{ "ErrorEquals": ["PagoRechazado"], "ResultPath": "$.error", "Next": "LiberarReserva" },
{ "ErrorEquals": ["States.ALL"], "ResultPath": "$.error", "Next": "LiberarReserva" }
],
"Next": "ConfirmarPedido"
},
"ConfirmarPedido": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": { "FunctionName": "km0-pedidos-cambiar-estado",
"Payload": { "pedido_id.$": "$.pedido_id", "estado": "pagado" } },
"ResultPath": null,
"Retry": [ { "ErrorEquals": ["States.ALL"], "IntervalSeconds": 1, "MaxAttempts": 5, "BackoffRate": 2.0 } ],
"Next": "Confirmado"
},
"Confirmado": { "Type": "Succeed" },
"LiberarReserva": {
"Type": "Task",
"Comment": "Compensación C2: idempotente (liberar dos veces no duplica stock)",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": { "FunctionName": "km0-inventario-liberar",
"Payload": { "id_reserva.$": "$.pedido_id" } },
"ResultPath": null,
"Retry": [ { "ErrorEquals": ["States.ALL"], "IntervalSeconds": 2, "MaxAttempts": 10, "BackoffRate": 2.0 } ],
"Next": "RechazarPedido"
},
"RechazarPedido": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": { "FunctionName": "km0-pedidos-cambiar-estado",
"Payload": { "pedido_id.$": "$.pedido_id", "estado": "rechazado", "motivo.$": "$.error.Error" } },
"ResultPath": null,
"Retry": [ { "ErrorEquals": ["States.ALL"], "IntervalSeconds": 2, "MaxAttempts": 10, "BackoffRate": 2.0 } ],
"Next": "Rechazado"
},
"Rechazado": { "Type": "Fail", "Error": "PedidoRechazado", "Cause": "Saga compensada" }
}
}Compárese con el OrquestadorSagaPedido de 03-05:
| Aspecto | saga_pedido.py + tabla sagas (03-05) |
Step Functions |
|---|---|---|
| Estado de la saga | Fila en sagas en km0_pedidos, escrita por nuestro código en cada transición |
Lo guarda el servicio; historial completo de cada ejecución consultable |
| Reintentos y backoff | Código propio (max_reintentos, FalloTransitorio) |
Declarativos por paso (Retry) |
| Compensaciones | Método _compensar que recorre los pasos hechos en orden inverso |
Ramas Catch → estados de compensación; el orden lo dibuja el grafo |
| Timeouts | Los del cliente gRPC (07-04) | TimeoutSeconds por paso, y de la ejecución completa |
| Recuperación tras caída del orquestador | Al arrancar, continuar() sobre las sagas en curso |
Automática: el servicio nunca "cae" para nosotros |
| Latencia por paso | Milisegundos (llamada gRPC directa) | Decenas de ms por transición + arranque de cada Lambda |
| Coste | El del servicio pedidos |
Por transición de estado (≈ 0,025 € por mil, tipo estándar): con 5 transiciones y 1,2 M pedidos/mes ≈ 150 €/mes |
| Pruebas | pytest con dobles (InventarioSimulado) y Testcontainers (07-06) |
Emulador local de Step Functions o staging; la lógica de cada paso, en Python |
| Portabilidad | Python, cualquier sitio | Solo AWS |
| Visibilidad | La que se instrumente (07-01/07-02) | Consola con el grafo y el estado de cada ejecución, de serie |
La decisión de Kilómetro Cero, que el ADR-006 de 08-05 formaliza: la saga de pedido se queda en pedidos, porque está en el camino crítico (latencia), tiene un volumen constante (coste) y ya está construida y probada; Step Functions se reserva para flujos de baja frecuencia y larga duración donde la visibilidad y los reintentos declarativos compensan, como la baja de un productor (cancelar sus productos, liquidar sus pagos, archivar sus fotos, todo con esperas de días).
- Edge computing: por qué acercar el cómputo
Todo lo anterior sigue en la región del proveedor. El edge computing mueve cómputo y datos hacia el borde: a los puntos de presencia de una CDN (cientos de ciudades), a las antenas o pasarelas del operador, o a los propios dispositivos. Cuatro razones, y para cada una un caso de Kilómetro Cero:
| Razón | Qué resuelve | Caso |
|---|---|---|
| Latencia | La luz tarda ~10 ms en 1 000 km de fibra (ida); una petición Valencia → Irlanda → Valencia son 40-60 ms antes de que el servidor haga nada | Las fotos y el catálogo de Ana desde un punto de presencia en Madrid o Valencia |
| Ancho de banda y egress | Cada foto servida desde la región paga egress y ocupa el enlace | La CDN sirve el 95 % de las fotos sin tocar S3 (04-03 lo calculaba) |
| Autonomía | Funcionar cuando la conexión con la nube falla | La app de furgoneta-3 en un túnel; el terminal de la Huerta La Vega en el mercado de Lleida sin cobertura |
| Datos locales | Datos que no hace falta (o no conviene) enviar enteros a la nube | Agregar 36 posiciones en una durante el túnel; filtrar en el dispositivo lo que no aporta |
- CDN: cachear fotos y catálogo cerca de Ana
Una CDN (Content Delivery Network: CloudFront, Cloudflare, Fastly, Akamai) es una red de servidores de caché distribuidos por el mundo, delante del origen. 04-05 la situó como capa de caché y 04-03 la justificó por el egress; aquí se diseña.
Cómo funciona. El DNS de km0.example apunta a la CDN; el navegador de Ana en Valencia resuelve al punto de presencia más cercano; si el objeto está en su caché (hit), lo sirve en 5-10 ms; si no (miss), lo pide al origen (S3 para las fotos, Kong para la API), lo guarda según las cabeceras de caché y lo sirve. El hit ratio es la métrica: con 4 200 productos y immutable, supera el 95 % tras el calentamiento.
Cache keys. La clave de caché es, por defecto, la URL completa (con parámetros de consulta). Dos consecuencias: miniatura-800.jpg?v=abc y ?v=def son objetos distintos, que es exactamente lo que se quiere para versionar (04-03); y una URL con parámetros irrelevantes (?utm_source=...) fragmenta la caché y baja el hit ratio, así que se configura qué parámetros forman parte de la clave. Cabeceras como Accept-Language o cookies pueden añadirse a la clave (con cuidado: cada valor multiplica las variantes).
TTL por tipo de contenido y invalidación:
| Contenido | Origen | Cache-Control |
Invalidación | Por qué |
|---|---|---|---|---|
| Miniaturas y originales de fotos | S3 km0-fotos |
public, max-age=31536000, immutable |
Nunca: la URL cambia con el contenido (?v=<etag>) |
Objetos inmutables (04-03) |
Ficha de producto (JSON de /api/v1/catalogo/productos/queso-curado) |
Kong → catalogo |
public, max-age=60, stale-while-revalidate=300 |
Por URL, disparada por el consumidor de stock.actualizado de 04-05 (además de Redis) |
Cambia poco; un minuto de obsolescencia es aceptable; stale-while-revalidate sirve lo viejo mientras refresca |
| Listado de un mercado | Kong → catalogo |
public, max-age=30 |
Por TTL | Cambia con cada publicación; barato de recalcular |
Cualquier cosa bajo /api/v1/pedidos |
Kong → pedidos |
private, no-store |
— | Datos personales; nunca en caché compartida |
Web estática (JS, CSS de seguimiento.js) |
S3 | immutable con hash en el nombre del fichero |
Nunca | Build con nombres hasheados |
WebSocket /ws |
reparto |
— | — | La CDN lo pasa (proxy) sin cachear; algunas terminan TLS y lo reenvían |
La invalidación explícita (purge por URL o por etiqueta) existe pero es lenta (segundos a minutos para propagarse a todos los puntos de presencia) y a menudo se cobra; el diseño correcto la evita para lo inmutable y la usa solo para lo semiestático (la ficha) como complemento del TTL corto. Es el mismo principio de 04-05: "TTL como red, invalidación explícita como mecanismo".
- Funciones en el borde: verificación de JWT y personalización
Las CDN modernas ejecutan código en el punto de presencia (Cloudflare Workers, Lambda@Edge y CloudFront Functions, Fastly Compute): funciones muy pequeñas (milisegundos, pocos MB) que ven la petición antes de que llegue al origen o la respuesta antes de que llegue al cliente. Sus usos en Kilómetro Cero:
- Verificación de JWT en el borde. Una petición con un token inválido o caducado se rechaza en Valencia, sin viajar a Irlanda ni ocupar a Kong: menos latencia para el error y menos carga en el origen ante un ataque. El borde verifica la firma (con la clave pública de Keycloak, cacheada) y la caducidad; la autorización fina sigue en el servicio (06-01), y Kong vuelve a verificar (defensa en profundidad, como en 08-02).
- Personalización sin perder la caché. La ficha de
queso-curadoes igual para todos salvo el mercado (precio y disponibilidad por ciudad). El Worker lee el mercado de una cookie o de la geolocalización de la petición y lo añade a la clave de caché (/productos/queso-curado?mercado=valencia), de modo que hay una variante por mercado y no una por usuario. - Redirecciones y A/B. Redirigir
/es/...segúnAccept-Language, enviar al 10 % de los usuarios a la nueva página de producto y medir, sin desplegar nada en el origen. - Servir desde el borde con respaldo. Si el origen falla, servir la última versión cacheada aunque haya caducado (
stale-if-error): el catálogo sigue visible durante un incidente.
// km0/borde/worker/catalogo.js — Cloudflare Worker: verifica el JWT y sirve del caché por mercado
import { jwtVerify, createRemoteJWKSet } from "jose"; // biblioteca JOSE; el JWKS de Keycloak se cachea en el borde
const JWKS = createRemoteJWKSet(new URL("https://id.km0.example/realms/km0/protocol/openid-connect/certs"));
export default {
async fetch(peticion, entorno, ctx) {
const url = new URL(peticion.url);
if (!url.pathname.startsWith("/api/v1/catalogo/")) return fetch(peticion); // el resto, al origen sin tocar
// 1. Verificación del JWT en el borde: firma, emisor, audiencia y caducidad (06-01)
const auth = peticion.headers.get("Authorization") || "";
const token = auth.startsWith("Bearer ") ? auth.slice(7) : null;
if (!token) return new Response(JSON.stringify({ error: "no_autenticado" }), { status: 401 });
let claims;
try {
({ payload: claims } = await jwtVerify(token, JWKS, { issuer: "https://id.km0.example/realms/km0", audience: "web-km0" }));
} catch (e) {
return new Response(JSON.stringify({ error: "token_invalido" }), { status: 401 });
}
// 2. Clave de caché por mercado (cookie o país de la petición), NO por usuario
const mercado = peticion.headers.get("Cookie")?.match(/mercado=([a-z]+)/)?.[1]
|| { ES: "valencia" }[peticion.cf?.country] || "girona";
const claveCache = new Request(`${url.origin}${url.pathname}?mercado=${mercado}`, { method: "GET" });
// 3. Servir del caché del punto de presencia si está
const cache = caches.default;
let respuesta = await cache.match(claveCache);
if (respuesta) return new Response(respuesta.body, { ...respuesta, headers: { ...Object.fromEntries(respuesta.headers), "X-Cache": "HIT" } });
// 4. Miss: ir al origen (Kong) con el mercado y el token (Kong vuelve a verificar: defensa en profundidad)
const origen = new Request(`${url.origin}${url.pathname}?mercado=${mercado}`, {
headers: { "Authorization": auth, "X-Request-Id": crypto.randomUUID(), "X-Mercado": mercado },
});
respuesta = await fetch(origen);
if (respuesta.ok && (respuesta.headers.get("Cache-Control") || "").includes("public")) {
ctx.waitUntil(cache.put(claveCache, respuesta.clone())); // guardar sin retrasar la respuesta
}
return respuesta;
},
};El Worker no contiene lógica de negocio ni accede a bases de datos: verifica, decide la clave de caché y reenvía. Todo lo que necesite el estado real (stock, pedidos) sigue yendo al origen. Y hay una restricción de diseño que el código respeta: la respuesta cacheada por mercado no puede contener datos del usuario; si catalogo añadiera "tus favoritos" a la ficha, la caché por mercado serviría los favoritos de Ana a Marc. Por eso la ficha es pública y los favoritos son otra llamada, private.
- Edge para IoT: la furgoneta y el mercado sin conexión
El borde más extremo es el dispositivo. Dos casos en Kilómetro Cero exigen funcionar sin conexión y sincronizar después:
La furgoneta. En 08-02, paho encolaba en memoria las posiciones durante el túnel y las enviaba de golpe al salir: 36 mensajes que a nadie interesaban uno a uno. Un agregador local en el dispositivo hace tres cosas mejor: guarda las posiciones en una cola persistente en disco (SQLite) para que un reinicio de la app no las pierda, agrega localmente (mientras no hay red, resume el tramo en un único mensaje con distancia recorrida, tiempo y trayectoria simplificada), y al recuperar la conexión envía primero la posición actual (lo urgente) y después el resumen (lo histórico). Es el edge como filtro: la nube recibe menos datos y más útiles.
El puesto del mercado. La Huerta La Vega vende en el mercado de Lleida con un terminal que descuenta stock local. Sin conexión, sigue vendiendo: es una réplica que acepta escrituras mientras está desconectada, es decir, replicación multilíder (03-04) con conflictos garantizados (la nube reserva 3 calabacines para un pedido online mientras el puesto vende 5 de los 6 que quedaban). 03-01 dio la herramienta para que la fusión sea determinista: los CRDT. Un contador de ventas del puesto como G-Counter (solo incrementa; fusionar es tomar el máximo por réplica) se sincroniza sin conflicto; el stock disponible, que es una resta, no es un CRDT puro, y se resuelve con la regla de negocio de 03-04: el puesto tiene autoridad sobre su stock físico (inv-lleida-puesto) y la nube sobre las reservas online, y la reconciliación puede dejar un pedido online sin stock, que la saga rechaza con compensación. Lo que el edge no puede hacer es prometer consistencia fuerte sin conexión: es la partición de 03-02, y el puesto elige disponibilidad.
# km0/borde/furgoneta/agregador.py
"""Agregador local en la furgoneta: cola persistente en SQLite, agregación sin conexión y
sincronización al volver. Complementa a furgoneta_mqtt.py (08-02)."""
import json, math, sqlite3, time
import paho.mqtt.client as mqtt
RUTA_DB = "/var/lib/km0/cola.sqlite"
T_POS = "km0/reparto/{id}/posicion"
T_TRAMO = "km0/reparto/{id}/tramo"
def _distancia_m(a, b) -> float:
"""Haversine simplificada entre (lat, lon) en metros."""
R = 6_371_000
p1, p2 = math.radians(a[0]), math.radians(b[0])
dp, dl = math.radians(b[0] - a[0]), math.radians(b[1] - a[1])
h = math.sin(dp / 2) ** 2 + math.cos(p1) * math.cos(p2) * math.sin(dl / 2) ** 2
return 2 * R * math.asin(math.sqrt(h))
class ColaPersistente:
"""Cola en disco: sobrevive a reinicios de la app y del dispositivo."""
def __init__(self, ruta: str = RUTA_DB):
self.db = sqlite3.connect(ruta)
self.db.execute("PRAGMA journal_mode=WAL") # escrituras rápidas y seguras
self.db.execute("""CREATE TABLE IF NOT EXISTS pendientes (
seq INTEGER PRIMARY KEY, fecha_ms INTEGER, lat REAL, lon REAL, enviada INTEGER DEFAULT 0)""")
def encolar(self, seq: int, fecha_ms: int, lat: float, lon: float) -> None:
with self.db:
self.db.execute("INSERT OR IGNORE INTO pendientes VALUES (?, ?, ?, ?, 0)", (seq, fecha_ms, lat, lon))
def pendientes(self) -> list[tuple]:
return self.db.execute("SELECT seq, fecha_ms, lat, lon FROM pendientes WHERE enviada = 0 ORDER BY seq").fetchall()
def marcar_enviadas(self, hasta_seq: int) -> None:
with self.db:
self.db.execute("UPDATE pendientes SET enviada = 1 WHERE seq <= ?", (hasta_seq,))
self.db.execute("DELETE FROM pendientes WHERE enviada = 1 AND fecha_ms < ?",
(int(time.time() * 1000) - 24 * 3600 * 1000,)) # limpiar lo de hace más de un día
class Agregador:
def __init__(self, repartidor: str, cliente: mqtt.Client, cola: ColaPersistente):
self.id, self.cliente, self.cola = repartidor, cliente, cola
self.conectado = False
cliente.on_connect = lambda *a: self._al_conectar()
cliente.on_disconnect = lambda *a: setattr(self, "conectado", False)
def _al_conectar(self) -> None:
self.conectado = True
self.sincronizar()
def nueva_posicion(self, seq: int, lat: float, lon: float) -> None:
fecha_ms = int(time.time() * 1000)
self.cola.encolar(seq, fecha_ms, lat, lon) # SIEMPRE a disco primero
if self.conectado:
self._publicar_posicion(seq, fecha_ms, lat, lon)
self.cola.marcar_enviadas(seq)
# sin conexión: no se hace nada más; la posición espera en disco
def _publicar_posicion(self, seq, fecha_ms, lat, lon) -> None:
carga = json.dumps({"repartidor": self.id, "seq": seq, "lat": lat, "lon": lon, "fecha_ms": fecha_ms})
self.cliente.publish(T_POS.format(id=self.id), carga, qos=1, retain=True)
def sincronizar(self) -> None:
"""Al recuperar la conexión: primero lo urgente (posición actual), luego el resumen del tramo."""
pendientes = self.cola.pendientes()
if not pendientes:
return
ultima = pendientes[-1]
self._publicar_posicion(*ultima) # 1. la posición actual, para el mapa de Ana
if len(pendientes) > 1: # 2. el tramo sin conexión, agregado en UN mensaje
puntos = [(p[2], p[3]) for p in pendientes]
distancia = sum(_distancia_m(puntos[i], puntos[i + 1]) for i in range(len(puntos) - 1))
tramo = {"repartidor": self.id, "seq_inicio": pendientes[0][0], "seq_fin": ultima[0],
"inicio_ms": pendientes[0][1], "fin_ms": ultima[1], "distancia_m": round(distancia),
"puntos": len(puntos), "trayectoria": puntos[::max(1, len(puntos) // 10)]} # 10 puntos como máximo
self.cliente.publish(T_TRAMO.format(id=self.id), json.dumps(tramo), qos=1)
self.cola.marcar_enviadas(ultima[0])Con el agregador, el ejercicio 2 de 08-02 cambia de respuesta: al salir del túnel, Ana recibe una posición (la actual), analitica recibe un tramo con la distancia y diez puntos (suficiente para el cálculo de distancia de Flink, que ahora puede sumar el distancia_m en lugar de reconstruirla), y el broker no recibe 36 mensajes. Los seq siguen siendo monótonos y persistidos, así que la deduplicación de 08-02 sigue funcionando. Lo que hay que añadir en el lado de la nube es un consumidor del nuevo tópico tramo en el puente MQTT → Kafka, hacia un tópico reparto.tramos.
- El continuo nube-edge-dispositivo y la seguridad en el borde
Nube, borde y dispositivo no son alternativas sino un continuo: cada cálculo y cada dato se sitúan donde el equilibrio entre latencia, ancho de banda, autonomía, capacidad y control sea mejor.
| Nivel | Qué hay | Capacidad | Latencia hacia el usuario | Autonomía | Qué pone Kilómetro Cero |
|---|---|---|---|---|---|
| Nube (región) | EKS, RDS, MSK, Cassandra, S3, Spark, Flink, Lambda | Ilimitada | 20-80 ms | Ninguna sin red | Todo lo transaccional y analítico; la verdad sobre pedidos, pagos, stock online |
| Borde de red (CDN, PoP) | Caché, Workers, terminación TLS, filtrado | Alta, pero funciones pequeñas y sin estado propio | 5-15 ms | Sirve caché si la nube falla | Fotos, catálogo, verificación de JWT, personalización por mercado, protección |
| Borde local (mercado, almacén) | Un terminal o mini-servidor | Baja, con disco | < 1 ms | Total, con sincronización posterior | Stock físico del puesto (inv-lleida-puesto), ventas presenciales, CRDT y reconciliación |
| Dispositivo (furgoneta, móvil) | App, cola SQLite, agregador | Mínima, batería | 0 | Total, con sincronización | Posiciones, agregación de tramos, comandos pendientes, mapa de Ana con última posición conocida |
flowchart TB
subgraph Nube[Nube: eu-west-1]
K[(Kafka)] --> S[Servicios en EKS] --> D[(RDS, Cassandra, S3)]
L[Lambda miniaturas, facturas]
end
subgraph Borde[Borde de red: CDN / PoP Valencia]
C[Caché de fotos y catálogo]
W[Worker: JWT, clave por mercado]
end
subgraph Local[Borde local: mercado de Lleida]
T[Terminal del puesto<br/>stock local, CRDT]
end
subgraph Disp[Dispositivos]
F[furgoneta-3<br/>cola SQLite, agregador]
A[Móvil de Ana<br/>última posición conocida]
end
A -- "HTTPS" --> W --> C -. "miss" .-> S
F -- "MQTT, QoS 1" --> K
T -. "sincronización periódica<br/>fusión CRDT" .-> S
S -- "WebSocket" --> A
Seguridad en el borde. Cuanto más lejos del centro, menos control físico: un Worker corre en infraestructura de un tercero, el terminal del mercado está en una mesa y la furgoneta puede ser robada. Cinco reglas:
- Nada secreto en el borde de red: el Worker verifica firmas con la clave pública; nunca tiene la privada ni credenciales de bases de datos. Los secretos que un Worker necesita (una clave de API del origen) se guardan en el almacén de secretos de la CDN, no en el código.
- El dispositivo se autentica con identidad propia y revocable: un certificado o credencial por dispositivo (la contraseña MQTT de
furgoneta-3de 08-02, emitida al darlo de alta), que se revoca en el broker si el dispositivo se pierde, y que solo autoriza su prefijo de tópicos (ACL). - Datos en reposo cifrados en el dispositivo: la cola SQLite con posiciones y los comandos con direcciones de clientes se cifran (SQLCipher o el almacén seguro del sistema operativo) y se borran cuando se sincronizan.
- Mínimo dato en el borde: el terminal del puesto no necesita la dirección ni el teléfono de Ana; recibe solo lo que le hace falta para vender (producto, unidades, precio), y lo que envía (ventas) no lleva datos personales.
- El borde no es la verdad: cualquier dato que venga del dispositivo se valida en la nube (el puente de 08-02 comprobaba que el
repartidordel payload coincide con el del tópico; eltramocon 900 km en 3 minutos se descarta como inválido). Un dispositivo comprometido puede mentir; el diseño debe limitar lo que una mentira puede causar.
Errores Comunes y Consejos
- Funciones para todo. Un consumidor constante de 500 eventos/s o un servicio con SLO de latencia son más baratos y previsibles como contenedor. Serverless para lo esporádico, lo evento-driven y lo que escala a cero.
- Suponer una sola ejecución. El proveedor reintenta; la función que no es idempotente genera dos facturas o tres miniaturas. Claves deterministas, marcas de idempotencia, operaciones condicionales.
- La función que se dispara a sí misma. Escribir la miniatura en el mismo bucket que dispara la función, sin filtro por prefijo o sufijo: bucle infinito y factura ilimitada. Filtro en el evento y comprobación en el código.
- Conexiones a base de datos desde mil instancias. Mil Lambdas abren mil conexiones a PostgreSQL. Límite de concurrencia, RDS Proxy, o DynamoDB para el estado de las funciones.
- Inicialización dentro del manejador. Cargar Pillow o crear el cliente de S3 en cada invocación multiplica la duración y el coste. Fuera del manejador, una vez por instancia.
- Sin DLQ. Un evento envenenado se reintenta y se pierde sin que nadie lo sepa. DLQ con alerta y procedimiento de reproceso, siempre.
- Cachear en la CDN respuestas con datos personales. Un
publicen/api/v1/pedidossirve los pedidos de Ana a Marc.private, no-storepara todo lo personal, y clave de caché por mercado, nunca por usuario, para lo compartido. - Sobrescribir objetos cacheados por la CDN. La caché no se entera. Claves inmutables con versión (
?v=<etag>), como se decidió en 04-03. - Confiar en el dispositivo. Valida en la nube todo lo que venga del borde; asume que puede estar comprometido.
- Consejo: separa siempre manejador (adaptador del proveedor) y lógica (Python puro). Las pruebas se hacen sobre la lógica; el lock-in se limita al adaptador.
- Consejo: calcula el coste de cada función a la frecuencia real y a diez veces la real. Si a diez veces sigue siendo barata, es un buen caso; si se dispara, prepara la salida a contenedor.
Ejercicios
Ejercicio 1: el punto de ack de la factura
En serverless/facturas/handler.py, reclamar(id_evento) se ejecuta antes de generar y escribir el PDF. (a) Describe el fallo concreto que esto produce y con qué probabilidad ocurre. (b) Reescribe la parte del manejador para que la idempotencia sea correcta, razonando por qué escribir dos veces en S3 es aceptable y qué ocurre si el proceso muere en cada punto posible. (c) ¿Sería distinta la respuesta si, en lugar de escribir un PDF en S3, la función enviara un correo con la factura?
Ejercicio 2: la CDN y la Semana del Queso Artesano
Durante la campaña, la web sirve 6 millones de vistas de producto al día; cada vista descarga una miniatura-800 (180 KB) y la ficha JSON (4 KB). (a) Con el diseño del apartado 9 (fotos immutable, ficha con max-age=60), estima el hit ratio esperable de cada tipo y el egress diario desde S3 y desde Kong, frente a servir todo desde el origen. (b) Un día, la Quesería Montblanc cambia la foto de queso-curado y se queja de que "algunos clientes siguen viendo la antigua". ¿Qué ha fallado, si el diseño es correcto? (c) El equipo propone añadir el nombre de usuario a la ficha ("Hola, Ana") para personalizarla. Explica el impacto en la caché y propón una alternativa.
Ejercicio 3: cuándo Step Functions
Para cada flujo, decide si Kilómetro Cero debería implementarlo con el orquestador propio de 03-05, con Step Functions, o con una simple función disparada por evento, y justifica con coste, latencia, duración y visibilidad: (a) la saga de confirmación de pedido (1,2 M al mes en campaña, en el camino crítico); (b) la baja de un productor (decenas al año: cancelar productos, esperar a que se entreguen los pedidos en curso, liquidar pagos a los 30 días, archivar fotos, notificar); (c) el reintento de cobros diferidos que la pasarela rechaza temporalmente (cientos al día, un reintento cada 6 horas durante 3 días).
Soluciones
Ejercicio 1.
(a) Si la función muere (timeout de 120 s por un lote grande, error de red con S3, retirada de la instancia) después de reclamar y antes de put_object, la marca factura#<id_evento> queda escrita y el PDF no. En el reintento de Lambda (o en el siguiente lote, porque el offset no avanzó), reclamar devuelve False y la factura se salta para siempre: factura perdida, sin error ni alerta. La probabilidad por evento es baja (la ventana son unos milisegundos entre dos llamadas), pero con 1,2 M pedidos al mes y reinicios de instancias, ocurrirá varias veces al mes.
(b) Reordenar: generar y escribir primero, marcar después; y usar la marca solo para evitar trabajo repetido, no para la corrección.
pedido = ev["datos"]
clave = f"facturas/{pedido['cliente']}/{pedido['pedido_id']}.pdf"
if ya_reclamado(ev["id_evento"]): # consulta, sin escribir: evita rehacer el PDF si ya se hizo
duplicados += 1; continue
s3.put_object(Bucket=BUCKET_FACTURAS, Key=clave, Body=generar_pdf(pedido), ...) # clave determinista
marcar(ev["id_evento"]) # PutItem sin condición (o con ella ignorando el fallo)Análisis por punto de muerte: antes de put_object: nada escrito, el reintento lo hace todo. Entre put_object y marcar: el PDF existe, la marca no; el reintento vuelve a generar el mismo PDF y lo escribe en la misma clave (S3 reemplaza el objeto entero; con versionado queda una versión más, inofensiva) y marca. Después de marcar: todo hecho. En ningún caso se pierde la factura, y el peor resultado es un PDF regenerado. Escribir dos veces es aceptable porque la operación es naturalmente idempotente: mismo contenido determinista (el PDF no lleva la hora de generación; si la llevara, habría que fijarla al fecha_ms del evento) en la misma clave.
(c) Sí. Enviar un correo no es idempotente por naturaleza: enviarlo dos veces molesta a Ana. Entonces la marca antes de enviar es necesaria para no duplicar, pero deja la ventana de pérdida. La solución es la de 02-05 con la pasarela: usar un servicio de correo que acepte una clave de idempotencia (muchos la ofrecen por mensaje), de modo que se pueda enviar "otra vez" sin duplicar; o dividir en dos pasos, generar el PDF (idempotente, en S3) y encolar el envío en una cola con deduplicación por id (SQS FIFO con MessageDeduplicationId), donde el envío se reintenta con seguridad. Sin una de esas dos cosas, hay que elegir entre perder o duplicar, y para un correo se prefiere duplicar.
Ejercicio 2.
(a) Fotos: 4 200 productos, objetos inmutables, un año de TTL; tras el calentamiento, cada punto de presencia tiene todas las miniaturas: hit ratio > 98 % (los misses son solo el primer acceso a cada objeto en cada PoP y las fotos nuevas). Egress desde S3: 6 M × 180 KB ≈ 1,08 TB/día sin CDN; con 98 % de hits, ≈ 22 GB/día desde S3 (más el egress de la CDN al usuario, que suele ser mucho más barato o incluido). Fichas: max-age=60 y 4 200 productos × 4 mercados = 16 800 variantes; 6 M vistas/día son 70 por segundo, repartidas en 16 800 claves: en promedio cada clave se pide cada 4 minutos, así que en la mayoría de PoP la ficha ha caducado cuando se vuelve a pedir; hit ratio quizá del 30-60 % (mejor para los productos populares, peor para la cola larga). stale-while-revalidate=300 sube ese ratio mucho (sirve la caducada y refresca en segundo plano), a costa de hasta 5 minutos de obsolescencia. Egress desde Kong: 6 M × 4 KB = 24 GB/día sin CDN; con 50 %, 12 GB. La foto domina: la CDN elimina el 98 % de ~1,1 TB diarios; para la ficha, el beneficio es de latencia y de carga en catalogo más que de egress.
(b) Si el diseño es correcto (la nueva foto tiene un ETag nuevo y la web referencia ?v=<etag nuevo>), lo que ha fallado es que la ficha que contiene la URL de la foto está cacheada: los clientes que reciben una ficha cacheada (hasta 60 s, o hasta 5 minutos con stale-while-revalidate) reciben la URL antigua, que sigue apuntando a un objeto inmutable y válido, la foto antigua. No es un fallo de la foto sino del TTL de la ficha, y es el comportamiento esperado: "algunos clientes durante unos minutos". Si el problema durara horas, la causa sería otra: la web referenciando la foto sin ?v= (sobrescritura de clave, el error de 04-03) o el consumidor de stock.actualizado invalidando Redis pero no la CDN (04-05 advertía de invalidar en todas las capas).
(c) Con "Hola, Ana" en el cuerpo de la ficha, la respuesta ya no es igual para todos los usuarios de un mercado: o se marca private (y se pierde la caché en la CDN: el hit ratio de la ficha cae a 0 y catalogo recibe los 70 req/s completos) o, peor, se cachea por mercado y Marc ve "Hola, Ana". La alternativa: la ficha sigue siendo pública y por mercado, y la personalización se hace en el cliente (el navegador ya tiene el nombre en el JWT o en la sesión y lo pinta) o con una segunda llamada pequeña y private (/api/v1/yo) que la CDN nunca cachea. Es el principio del Worker: separar lo compartido cacheable de lo personal.
Ejercicio 3.
(a) Orquestador propio en pedidos. Volumen alto y constante (≈ 150 €/mes solo en transiciones de Step Functions, más las invocaciones), en el camino crítico (cada transición añade decenas de ms al p99, comprometiendo el SLO de 500 ms), duración de segundos, y ya construido y probado con Testcontainers. La visibilidad la dan las trazas de 07-02 y la tabla sagas.
(b) Step Functions. Decenas al año (coste despreciable), duración de semanas (esperas de 30 días que un orquestador propio tendría que gestionar con estado persistido y temporizadores: Step Functions tiene el estado Wait con fechas y ejecuciones de hasta un año), fuera de cualquier camino crítico, y con enorme valor de visibilidad (ver en qué paso está la baja de cada productor, quién la aprobó, por qué falló la liquidación). Cada paso invoca las APIs de los servicios existentes. El lock-in se acepta porque el flujo es periférico.
(c) Una función disparada por temporizador (o una cola con retraso), no Step Functions ni orquestador. Cientos al día con un reintento cada 6 horas durante 3 días es un patrón de cola con reintento diferido: SQS con DelaySeconds (máximo 15 minutos, así que se encadena) o, más simple, una tabla cobros_pendientes con siguiente_intento y una Lambda programada cada 15 minutos que procesa los vencidos con la Idempotency-Key de 02-05 y publica pago.confirmado o pago.rechazado al terminar. Step Functions con Wait de 6 horas funcionaría, pero es una máquina de estados para lo que es un bucle con una fecha; el coste y la complejidad no se justifican, y la visibilidad la da la tabla.
Conclusión
Serverless y edge computing amplían el espacio de decisiones de Kilómetro Cero en dos direcciones opuestas a la del clúster en una región. Con FaaS, el proveedor ejecuta funciones efímeras en respuesta a eventos, escala de cero a miles y cobra por invocación y GB-segundo; a cambio impone arranques en frío, límites de tiempo y memoria, ausencia de estado y, sobre todo, reintentos que no controlamos y que hacen de la idempotencia de 02-05 una obligación y no una buena práctica. Con ese filtro, las miniaturas de km0-fotos, las facturas de pago.confirmado, los webhooks de la pasarela y las tareas programadas salen de los servicios y pasan a funciones definidas con SAM, con DLQ, límite de concurrencia y filtros que evitan bucles, mientras pedidos, el WebSocket de 08-02, Flink y Spark se quedan donde estaban. Los servicios serverless de alrededor (colas, DynamoDB, API Gateway, Step Functions) completan el modelo, y la saga de 03-05 reescrita en Amazon States Language muestra con precisión qué se gana (estado gestionado, reintentos declarativos, visibilidad) y qué se paga (latencia, coste por transición, lock-in), lo que deja la saga de pedido en pedidos y reserva Step Functions para los flujos largos y esporádicos. En la otra dirección, el edge acerca datos y cómputo a quien los usa: la CDN sirve fotos inmutables y fichas por mercado con TTL e invalidación diseñados por tipo de contenido, el Worker verifica el JWT y decide la clave de caché sin tocar el origen, el agregador de la furgoneta convierte 36 posiciones de túnel en una posición y un tramo, y el puesto del mercado sigue vendiendo sin conexión con CRDT y reconciliación por autoridad, aceptando la partición de 03-02 en lugar de negarla. El continuo nube-borde-dispositivo sitúa cada dato donde el equilibrio es mejor, y la seguridad en el borde parte de que el borde puede mentir.
Con esto, todas las piezas de Kilómetro Cero están sobre la mesa: desde el corte de los bounded contexts de 08-01 hasta la caché en Valencia y la cola SQLite en la furgoneta. La última lección las junta: la arquitectura completa en un solo diagrama, el recorrido de un pedido de Ana de extremo a extremo con todo lo que deja a su paso, las decisiones y sus trade-offs revisitados como ADR, los cinco síntomas de 01-06 con su resolución, la evaluación con los criterios del curso, lo que queda fuera, el proyecto que el alumno puede construir en su portátil y la guía para seguir profundizando. Es el Proyecto Final: Kilómetro Cero de Extremo a Extremo.
Curso de Arquitecturas Distribuidas
Módulo 1: Introducción a los Sistemas Distribuidos
- Conceptos Básicos de Sistemas Distribuidos
- Modelos de Sistemas Distribuidos
- Ventajas y Desafíos de los Sistemas Distribuidos
- Las Falacias de la Computación Distribuida
- Tiempo, Relojes y Ordenación de Eventos
- Del Monolito a la Plataforma Distribuida: el Caso Kilómetro Cero
Módulo 2: Comunicación en Sistemas Distribuidos
- Protocolos de Comunicación
- RPC y RMI
- gRPC y Serialización de Datos
- Mensajería y Colas de Mensajes
- Patrones de Comunicación Asíncrona
Módulo 3: Consistencia y Replicación
- Modelos de Consistencia
- El Teorema CAP y PACELC
- Algoritmos de Consenso
- Replicación de Datos
- Transacciones Distribuidas y Sagas
Módulo 4: Almacenamiento Distribuido
- Particionado de Datos y Hashing Consistente
- Sistemas de Archivos Distribuidos
- Almacenamiento de Objetos
- Bases de Datos Distribuidas
- Cachés Distribuidos
Módulo 5: Computación Distribuida
- Modelos de Computación Distribuida
- MapReduce y Hadoop
- Spark y Computación en Memoria
- Procesamiento de Flujos de Datos
- Planificación de Trabajos y Pipelines de Datos
Módulo 6: Seguridad en Sistemas Distribuidos
- Autenticación y Autorización
- Cifrado y Protección de Datos
- Gestión de Identidades
- Seguridad entre Servicios: mTLS y Gestión de Secretos
- Puertas de Enlace, Limitación de Tasa y Auditoría
Módulo 7: Monitoreo y Mantenimiento
- Monitoreo de Sistemas Distribuidos
- Logs Centralizados y Trazabilidad Distribuida
- Gestión de Fallos y Recuperación
- Patrones de Resiliencia: Timeouts, Reintentos y Circuit Breaker
- Automatización y Orquestación
- Pruebas en Sistemas Distribuidos e Ingeniería del Caos
