La lección anterior dejó a Jordi Sala con una alerta en la mano: PedidosBurnRateRapido, el 8 % de los pedidos falla. Las métricas le dicen cuánto y desde cuándo; no le dicen qué pedidos, ni en qué servicio, ni por qué. En el monolito de Kilómetro Cero la respuesta estaba a un grep de distancia en un único fichero de log. Hoy un pedido atraviesa Kong, tres réplicas de pedidos, dos de inventario, pagos, un tópico de Kafka y analitica, cada uno escribiendo en su propio contenedor; el log del pedido P-2026-000125 está repartido en siete sitios y en ninguno se llama igual. Esta lección construye los otros dos pilares de la observabilidad: logs estructurados y centralizados, para que un identificador reúna todo lo que pasó con una petición, y trazas distribuidas, para ver el recorrido y el tiempo de esa petición por todos los servicios. Al final, las tres señales quedan unidas en Grafana: de la métrica que alerta, a la traza que localiza, al log que explica.
Contenido
- Por qué
tail -fya no vale - Logs estructurados: campos estándar, niveles y qué no loguear
- Pipeline de recogida: agente, almacén y consulta
servicios/comun/logs.py: JSON contrace_idinyectado- Trazabilidad distribuida: trazas, spans y contexto
- Instrumentación con OpenTelemetry:
servicios/comun/trazas.py - Collector, backends y muestreo
- Una traza de "crear pedido" comentada
- Unir las tres señales
- Errores comunes y consejos
- Ejercicios y soluciones
- Conclusión
- Por qué
tail -f ya no vale
tail -f ya no valeEn una máquina, un log es un fichero y un incidente se investiga con tail -f, grep y less. En un sistema distribuido eso deja de funcionar por razones que no son de comodidad, sino estructurales:
- La información de una petición está repartida. Crear el pedido
P-2026-000125genera líneas en Kong,pedidos,inventario,pagosyanalitica. Para reconstruir la historia hay que abrir cinco terminales y adivinar qué línea de cada una corresponde a ese pedido. - Los relojes no coinciden. Como se vio en 01-05, ordenar por timestamp líneas de máquinas distintas puede invertir causa y efecto. Hace falta un identificador de correlación, no solo la hora.
- Los contenedores son efímeros. Cuando Kubernetes (07-05) reemplaza un pod de
pedidos, su log desaparece con él. Si el fallo ocurrió justo antes de morir, se ha perdido la evidencia. - El formato libre no se consulta.
Pedido 125 rechazado por stock (tomate-rosa)es legible para una persona; para preguntar "cuántos rechazos por stock detomate-rosaen Lleida en la última hora" hay que parsear con expresiones regulares frágiles. - Volumen. 140 repartidores enviando una posición cada 5 s, 100 peticiones/s en
pedidos: son millones de líneas al día. Nadie las lee; hay que poder filtrarlas.
La solución tiene dos partes: que cada servicio escriba logs estructurados con campos comunes, y que un pipeline los recoja de todos los nodos en un almacén central consultable.
- Logs estructurados: campos estándar, niveles y qué no loguear
Un log estructurado es un registro con campos con nombre, normalmente una línea JSON por evento:
{"ts": "2026-09-12T10:13:41.208Z", "nivel": "warning", "servicio": "pedidos", "instancia": "pedidos-2",
"trace_id": "4bf92f3577b34da6a3ce929d0e0e4736", "span_id": "00f067aa0ba902b7",
"request_id": "c1d2e3f4-5a6b-7c8d-9e0f-1a2b3c4d5e6f", "usuario": "u:9f1c3a",
"evento": "reserva_stock_rechazada", "pedido_id": "P-2026-000125",
"producto": "vino-crianza", "solicitadas": 6, "disponibles": 2, "replica": "inv-bcn"}Los campos que todos los servicios de km0/ incluyen siempre:
| Campo | Contenido | Por qué |
|---|---|---|
ts |
Instante en UTC, ISO 8601 con milisegundos | Sin zona horaria no hay ordenación posible entre nodos; UTC evita el cambio de hora |
nivel |
debug, info, warning, error, critical |
Filtrar y alertar |
servicio |
pedidos, inventario, ... |
Primer filtro en cualquier consulta |
instancia |
pedidos-2, nombre del pod |
Distinguir una réplica enferma |
trace_id |
Identificador de la traza OpenTelemetry (32 hex) | Correlación entre servicios y con las trazas |
span_id |
Span actual (16 hex) | Enlazar la línea con el paso concreto de la traza |
request_id |
El X-Request-Id que Kong genera y propaga (06-05) |
Es el identificador que ve el cliente y que Ana puede dar a soporte |
usuario |
Identificador seudonimizado del sujeto (06-02) | Investigar sin exponer identidad |
evento |
Nombre corto y estable del suceso (pedido_creado, reserva_stock_rechazada) |
Contar y filtrar por tipo sin parsear el mensaje |
mensaje (opcional) |
Texto para humanos | Contexto adicional |
Los campos específicos (pedido_id, producto, replica) se añaden como claves adicionales, no dentro del mensaje. Y sí: aquí pedido_id es bienvenido. Lo que en las métricas era un problema de cardinalidad (07-01) en los logs es exactamente lo que se busca, porque un log se indexa por texto o por unas pocas etiquetas, no crea una serie por valor.
Niveles
| Nivel | Cuándo | Ejemplo en pedidos |
|---|---|---|
debug |
Detalle para desarrollo; apagado en producción salvo investigación puntual | Cuerpo de la petición gRPC a inventario |
info |
Hitos normales del negocio | pedido_creado, saga_completada |
warning |
Algo inesperado que el sistema ha manejado | reserva_stock_rechazada, reintento de una llamada |
error |
Una operación ha fallado y alguien debería mirarlo | pago_error_tecnico, compensacion_fallida |
critical |
El servicio no puede seguir | No conecta con km0_pedidos al arrancar |
Qué no loguear
Los logs viajan por la red, se almacenan meses y los lee mucha gente. Todo lo que el Módulo 6 protegió con cifrado y control de acceso puede acabar en claro en un log si no se vigila:
- Secretos: tokens JWT (ni siquiera "para depurar"), contraseñas, credenciales de Vault, claves de API de la pasarela de pagos, cabeceras
Authorizationcompletas. Una traza HTTP endebugque vuelque cabeceras es una fuga. - Datos personales en claro: nombre, email, teléfono y dirección de Ana. Se registra
usuario: "u:9f1c3a"(el seudónimo de 06-02) ypedido_id; quien necesite el teléfono lo obtiene de la base de datos con su permiso y su auditoría. - Datos de tarjeta: nunca, en ninguna forma (06-02).
- Cuerpos completos de peticiones: además de datos personales, son enormes. Se loguean los campos relevantes escogidos uno a uno.
Los logs de auditoría de 06-05 (quién hizo qué, hash encadenado, object lock en MinIO) son un canal distinto con garantías distintas: inmutabilidad y retención larga. Los logs operativos de esta lección son para diagnosticar, se retienen semanas y se pueden borrar; no se deben mezclar.
- Pipeline de recogida: agente, almacén y consulta
flowchart LR
subgraph nodo1["Nodo 1"]
P1[pedidos-1<br/>stdout JSON]
I1[inventario-1<br/>stdout JSON]
A1[Promtail]
P1 --> A1
I1 --> A1
end
subgraph nodo2["Nodo 2"]
P2[pedidos-2]
K[Kong]
A2[Promtail]
P2 --> A2
K --> A2
end
L[(Loki<br/>índice por etiquetas<br/>chunks en MinIO km0-logs)]
G[Grafana<br/>LogQL]
A1 --> L
A2 --> L
G --> L
Cada servicio escribe JSON en stdout (no en ficheros propios: el contenedor no debe saber dónde acaban sus logs). Un agente por nodo lee la salida de todos los contenedores del nodo, añade etiquetas (servicio, instancia, nodo) y envía al almacén central. La consulta se hace desde una interfaz que entiende el formato.
| Componente | Opción "Loki" | Opción "ELK" | Otras |
|---|---|---|---|
| Agente | Promtail (o Grafana Alloy) | Filebeat, Logstash | Fluent Bit / Fluentd (ambas pilas), Vector |
| Almacén | Loki | Elasticsearch / OpenSearch | — |
| Consulta | Grafana (LogQL) | Kibana / OpenSearch Dashboards (KQL, Lucene) | — |
Loki frente a Elasticsearch, que es la decisión de arquitectura real:
| Aspecto | Loki | Elasticsearch |
|---|---|---|
| Qué indexa | Solo las etiquetas (servicio, instancia, nivel); el contenido se guarda comprimido sin indexar |
El texto completo de cada campo (índice invertido) |
| Coste de almacenamiento | Bajo (chunks comprimidos en un bucket de objetos, p. ej. MinIO 04-03) | Alto (el índice puede ocupar más que los datos) |
| Coste de ingesta | Bajo | Alto en CPU y memoria |
Consulta "todo lo de pedidos con trace_id=X en la última hora" |
Filtra por etiqueta y luego escanea el contenido de esa hora: rápido si el rango es acotado | Instantáneo por el índice |
| Búsqueda libre en meses de datos | Lenta (escaneo) | Rápida |
| Analítica sobre logs (agregaciones complejas) | Limitada (LogQL tiene métricas sobre logs, pero no es un motor analítico) | Potente |
| Integración | Nativa con Grafana, Prometheus y Tempo: mismo lenguaje de etiquetas, enlaces log↔traza↔métrica | Kibana; APM propio |
| Cuándo | Diagnóstico operativo con etiquetas conocidas, presupuesto contenido, pila Grafana | Búsqueda exploratoria intensa, cumplimiento con búsquedas ad hoc, equipos que ya operan Elasticsearch |
Kilómetro Cero elige Loki porque su patrón de consulta es casi siempre "servicio + rango + trace_id", el almacenamiento va a MinIO y ya usa Grafana y Prometheus. La cardinalidad vuelve a importar, esta vez en las etiquetas de Loki: servicio, instancia y nivel son etiquetas; trace_id, pedido_id y usuario no lo son (crearían un stream por valor), se buscan dentro del contenido con | json | trace_id="...".
Retención y coste
Los logs crecen sin límite si nadie decide lo contrario. Una política típica: debug no se envía en producción; info se conserva 14 días; warning/error 90 días; los logs de Kong (una línea por petición) 30 días; los de auditoría son otro sistema con años de retención. En Loki la retención se configura por stream y los chunks antiguos se borran del bucket. Reducir el volumen en origen (no loguear cada posición de reparto, sino un resumen por minuto y los rechazos) suele ser la medida más eficaz.
Configuración mínima de Promtail y Loki en docker-compose.yml:
# km0/docker-compose.yml (fragmento)
loki:
image: grafana/loki:3.1.0
command: ["-config.file=/etc/loki/loki.yaml"]
volumes:
- ./observabilidad/loki.yaml:/etc/loki/loki.yaml:ro
ports: ["3100:3100"]
promtail:
image: grafana/promtail:3.1.0
command: ["-config.file=/etc/promtail/promtail.yaml"]
volumes:
- ./observabilidad/promtail.yaml:/etc/promtail/promtail.yaml:ro
- /var/lib/docker/containers:/var/lib/docker/containers:ro
- /var/run/docker.sock:/var/run/docker.sock:ro# km0/observabilidad/promtail.yaml
server:
http_listen_port: 9080
clients:
- url: http://loki:3100/loki/api/v1/push
scrape_configs:
- job_name: contenedores
docker_sd_configs:
- host: unix:///var/run/docker.sock
relabel_configs:
# La etiqueta de compose "com.docker.compose.service" pasa a ser la etiqueta "servicio"
- source_labels: [__meta_docker_container_label_com_docker_compose_service]
target_label: servicio
- source_labels: [__meta_docker_container_name]
regex: "/(.*)"
target_label: instancia
pipeline_stages:
- json:
expressions:
nivel: nivel
- labels:
nivel: # solo "nivel" se promociona a etiqueta; trace_id se queda en el contenido# km0/observabilidad/loki.yaml (fragmento relevante)
storage_config:
aws:
s3: s3://km0-logs
endpoint: minio:9000
access_key_id: ${LOKI_MINIO_KEY}
secret_access_key: ${LOKI_MINIO_SECRET}
s3forcepathstyle: true
limits_config:
retention_period: 336h # 14 días por defecto
compactor:
retention_enabled: true
servicios/comun/logs.py: JSON con trace_id inyectado
servicios/comun/logs.py: JSON con trace_id inyectadoCon structlog, cada servicio configura una vez el logger y después usa log.info("evento", clave=valor). Un processor consulta el contexto de OpenTelemetry (sección 6) y añade trace_id y span_id si hay un span activo:
# km0/servicios/comun/logs.py
"""Logs JSON estructurados con correlación por trace_id / request_id."""
import logging
import os
import sys
from contextvars import ContextVar
import structlog
from opentelemetry import trace
# El request_id de Kong se guarda en una ContextVar: es local a la petición
# en curso (funciona con hilos y con asyncio) y cualquier log lo puede leer.
request_id_actual: ContextVar[str | None] = ContextVar("request_id", default=None)
SERVICIO = os.environ.get("KM0_SERVICIO", "desconocido")
INSTANCIA = os.environ.get("HOSTNAME", "local")
# Claves que jamás deben salir en un log, aunque alguien las pase por descuido.
CLAVES_PROHIBIDAS = {"authorization", "password", "contrasena", "token", "tarjeta", "cvv", "telefono", "email"}
def inyectar_contexto(logger, metodo, evento: dict) -> dict:
"""Processor de structlog: añade servicio, instancia, trace_id, span_id y request_id."""
evento["servicio"] = SERVICIO
evento["instancia"] = INSTANCIA
span = trace.get_current_span()
ctx = span.get_span_context()
if ctx.is_valid:
# format(x, "032x") = 32 dígitos hexadecimales, el formato W3C
evento["trace_id"] = format(ctx.trace_id, "032x")
evento["span_id"] = format(ctx.span_id, "016x")
rid = request_id_actual.get()
if rid:
evento["request_id"] = rid
return evento
def censurar(logger, metodo, evento: dict) -> dict:
"""Processor: sustituye el valor de cualquier clave sensible por '[censurado]'."""
for clave in list(evento):
if clave.lower() in CLAVES_PROHIBIDAS:
evento[clave] = "[censurado]"
return evento
def configurar_logs(nivel: str = "INFO") -> None:
structlog.configure(
processors=[
structlog.contextvars.merge_contextvars, # campos ligados con bind_contextvars
structlog.processors.add_log_level, # -> "level"
structlog.processors.TimeStamper(fmt="iso", utc=True, key="ts"),
inyectar_contexto,
censurar,
structlog.processors.EventRenamer("evento"), # el primer argumento pasa a "evento"
structlog.processors.JSONRenderer(),
],
wrapper_class=structlog.make_filtering_bound_logger(getattr(logging, nivel)),
logger_factory=structlog.PrintLoggerFactory(file=sys.stdout),
)
log = structlog.get_logger()Cómo se usa en pedidos, con el middleware que captura el X-Request-Id que Kong propaga:
# km0/servicios/pedidos/app.py (fragmento)
from servicios.comun.logs import log, request_id_actual, configurar_logs
configurar_logs(nivel=os.environ.get("KM0_LOG_NIVEL", "INFO"))
@app.middleware("http")
async def middleware_request_id(request: Request, call_next):
rid = request.headers.get("x-request-id", "sin-request-id")
token = request_id_actual.set(rid)
try:
return await call_next(request)
finally:
request_id_actual.reset(token)
# En la saga (03-05), al rechazarse una reserva:
log.warning("reserva_stock_rechazada", pedido_id=pedido.id, producto=linea.producto,
solicitadas=linea.cantidad, disponibles=resp.disponibles, replica=resp.replica)Detalles que conviene notar:
structlogno formatea cadenas:log.warning("reserva_stock_rechazada", producto=...)produce claves separadas. Escribirlog.warning(f"Rechazado {producto}")funciona, pero pierde la estructura. La disciplina es "nombre de evento estable + campos".- El processor
censurares una red de seguridad, no la política: la política es no pasar esos datos. Pero cuesta poco y evita que unlog.debug("peticion", headers=dict(request.headers))de un martes por la noche acabe volcando tokens. - El
trace_idse lee del span activo de OpenTelemetry. Sin trazas configuradas, el contexto es inválido y simplemente no aparece; con ellas, cada línea de log de cualquier servicio que participe en la misma petición lleva el mismotrace_id, aunque los servicios no hayan hecho nada explícito para pasarlo.
Consultas LogQL en Grafana:
# Todo lo que pasó con una traza concreta, en todos los servicios, en orden
{servicio=~"pedidos|inventario|pagos|analitica"} | json | trace_id="4bf92f3577b34da6a3ce929d0e0e4736"
# Rechazos de stock por producto en la última hora (métrica derivada de logs)
sum by (producto) (count_over_time({servicio="pedidos"} | json | evento="reserva_stock_rechazada" [1h]))
# Errores de pedidos-2 que no son de la pasarela de pagos
{servicio="pedidos", instancia="pedidos-2", nivel="error"} | json | evento != "pago_error_tecnico"
# Lo que Ana reporta a soporte: su request_id
{servicio=~".+"} | json | request_id="c1d2e3f4-5a6b-7c8d-9e0f-1a2b3c4d5e6f"La segunda consulta muestra un uso valioso: métricas a partir de logs para preguntas que no merecen una métrica propia (07-01, ejercicio 1) pero que sí se quieren graficar de vez en cuando.
- Trazabilidad distribuida: trazas, spans y contexto
Los logs con trace_id responden "qué pasó"; las trazas responden "por dónde pasó y cuánto tardó cada paso". El modelo, heredado de Dapper (Google) y estandarizado por OpenTelemetry:
- Una traza es el árbol completo de trabajo causado por una petición externa. Se identifica por un
trace_idde 128 bits. - Un span es una unidad de trabajo con nombre, inicio, fin, atributos (
pedido_id,rpc.method,db.statement), estado (OK/ERROR) y eventos. Tienespan_idy unparent_span_id, salvo el raíz. - El contexto de propagación es lo que viaja de un servicio a otro para que el span del receptor cuelgue del emisor. El estándar es W3C Trace Context: una cabecera
traceparent:
traceparent: 00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01
^^ ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ ^^^^^^^^^^^^^^^^ ^^
versión trace_id (32 hex) span padre flags (01 = muestreada)(Y una opcional tracestate para datos de proveedores.) En HTTP y gRPC es una cabecera; en Kafka, una cabecera del mensaje; en la saga, se guarda junto al estado en la tabla sagas para que la compensación horas después siga colgando de la traza original.
sequenceDiagram
participant Ana
participant Kong
participant Ped as pedidos-2
participant Inv as inventario-1
participant Pag as pagos
participant Kafka
participant Ana2 as analitica
Ana->>Kong: POST /api/v1/pedidos
Note over Kong: crea trace_id 4bf9...<br/>span raíz "kong.proxy"
Kong->>Ped: POST /pedidos<br/>traceparent: 00-4bf9...-a1..-01<br/>X-Request-Id: c1d2...
Note over Ped: span "POST /pedidos" hijo de a1
Ped->>Inv: gRPC ReservarStock<br/>metadata traceparent: 00-4bf9...-b2..-01
Note over Inv: span "ReservarStock" + span "SELECT ... FOR UPDATE"
Inv-->>Ped: OK
Ped->>Pag: gRPC Cobrar (traceparent ...-b2..)
Pag-->>Ped: OK
Ped->>Kafka: pedido.confirmado<br/>header traceparent: 00-4bf9...-c3..-01
Ped-->>Kong: 201
Kong-->>Ana: 201
Kafka-->>Ana2: consume (140 ms después)
Note over Ana2: span "consumir pedidos.eventos"<br/>enlazado a c3
- Instrumentación con OpenTelemetry:
servicios/comun/trazas.py
servicios/comun/trazas.pyOpenTelemetry (OTel) aporta tres cosas: un SDK que crea spans y gestiona el contexto, instrumentaciones automáticas para librerías conocidas (gRPC, FastAPI, psycopg, redis, requests, kafka-python/confluent-kafka), y un protocolo de exportación (OTLP) hacia un Collector o un backend. Los servicios de km0/ comparten la configuración:
# km0/servicios/comun/trazas.py
"""Configuración de OpenTelemetry para trazas distribuidas."""
import os
from opentelemetry import trace, context, propagate
from opentelemetry.sdk.resources import Resource, SERVICE_NAME, SERVICE_INSTANCE_ID
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from opentelemetry.sdk.trace.sampling import ParentBased, TraceIdRatioBased
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.instrumentation.grpc import GrpcInstrumentorServer, GrpcInstrumentorClient
from opentelemetry.instrumentation.psycopg import PsycopgInstrumentor
from opentelemetry.instrumentation.redis import RedisInstrumentor
from opentelemetry.propagators.textmap import Getter, Setter
def configurar_trazas(servicio: str) -> trace.Tracer:
recurso = Resource.create({
SERVICE_NAME: servicio, # "pedidos", "inventario"...
SERVICE_INSTANCE_ID: os.environ.get("HOSTNAME", "local"),
"deployment.environment": os.environ.get("KM0_ENTORNO", "produccion"),
})
# ParentBased: si la petición llega ya con decisión de muestreo (flags=01), se respeta;
# si es raíz, se muestrea el 10 %. Así una traza nunca queda "a medias".
muestreo = ParentBased(root=TraceIdRatioBased(float(os.environ.get("KM0_TRAZAS_RATIO", "0.1"))))
proveedor = TracerProvider(resource=recurso, sampler=muestreo)
exportador = OTLPSpanExporter(endpoint=os.environ.get("OTEL_EXPORTER_OTLP_ENDPOINT", "otel-collector:4317"),
insecure=False) # mTLS con el certificado de Vault (06-04)
proveedor.add_span_processor(BatchSpanProcessor(exportador)) # exporta en lotes, en un hilo aparte
trace.set_tracer_provider(proveedor)
# Auto-instrumentación: parchean las librerías para crear spans y propagar traceparent
GrpcInstrumentorServer().instrument()
GrpcInstrumentorClient().instrument()
PsycopgInstrumentor().instrument(enable_commenter=True) # añade el trace_id como comentario SQL
RedisInstrumentor().instrument()
return trace.get_tracer(servicio)
# ---- Propagación manual por cabeceras Kafka -----------------------------------
# Las cabeceras Kafka son una lista de (clave, bytes); OTel necesita un Getter/Setter
# que sepa leerlas y escribirlas.
class _KafkaSetter(Setter):
def set(self, carrier: list, key: str, value: str) -> None:
carrier.append((key, value.encode("utf-8")))
class _KafkaGetter(Getter):
def get(self, carrier: list, key: str):
return [v.decode("utf-8") for k, v in carrier if k == key] or None
def keys(self, carrier: list):
return [k for k, _ in carrier]
def inyectar_en_cabeceras_kafka() -> list[tuple[str, bytes]]:
"""Devuelve cabeceras Kafka con el traceparent del span activo."""
cabeceras: list[tuple[str, bytes]] = []
propagate.inject(cabeceras, setter=_KafkaSetter())
return cabeceras
def contexto_desde_cabeceras_kafka(cabeceras: list[tuple[str, bytes]] | None):
"""Reconstruye el contexto de traza a partir de las cabeceras de un mensaje."""
return propagate.extract(cabeceras or [], getter=_KafkaGetter())Cómo lo usa pedidos al publicar el evento tras la saga, y analitica al consumirlo:
# km0/servicios/pedidos/eventos.py (fragmento)
from servicios.comun.trazas import inyectar_en_cabeceras_kafka
def publicar_pedido_confirmado(productor, pedido):
with tracer.start_as_current_span("publicar pedidos.eventos",
kind=trace.SpanKind.PRODUCER,
attributes={"messaging.system": "kafka",
"messaging.destination": "pedidos.eventos",
"km0.pedido_id": pedido.id}):
productor.produce("pedidos.eventos", key=pedido.id.encode(), value=pedido.a_json().encode(),
headers=inyectar_en_cabeceras_kafka()) # traceparent viaja con el mensaje# km0/servicios/analitica/consumidor_pedidos.py (fragmento)
from opentelemetry import context, trace
from servicios.comun.trazas import contexto_desde_cabeceras_kafka
for msg in consumidor:
ctx = contexto_desde_cabeceras_kafka(msg.headers())
# El span del consumidor es hijo del span "publicar" de pedidos: misma traza, 140 ms después
with tracer.start_as_current_span("consumir pedidos.eventos", context=ctx,
kind=trace.SpanKind.CONSUMER,
attributes={"messaging.kafka.partition": msg.partition(),
"messaging.kafka.offset": msg.offset()}):
procesar(msg)
log.info("evento_procesado", pedido_id=msg.key().decode()) # lleva el trace_id de pedidosDos aclaraciones:
- gRPC, HTTP y PostgreSQL no requieren código: los instrumentors interceptan las llamadas, crean spans y meten/leen
traceparenten los metadatos. Eninventario,AuthInterceptory elMetricasInterceptorde 07-01 conviven con el interceptor de OTel; el orden no afecta al contexto. - Kafka sí requiere propagación manual (o el instrumentor de
confluent-kafka, que hace lo mismo por debajo): el mensaje es un dato inerte que se consume más tarde, en otro proceso; nadie lo "llama". Lo mismo ocurre con la saga:saga_pedido.pyguardatraceparenten la fila de la tablasagas, y cuando el relay del outbox o una compensación retoman el trabajo, reconstruyen el contexto conpropagate.extractsobre ese valor. Sin esto, la compensación deP-2026-000125aparecería como una traza huérfana sin relación con el pedido.
- Collector, backends y muestreo
Los servicios no envían las trazas al backend directamente, sino al OpenTelemetry Collector: un proceso intermedio que recibe OTLP, procesa (añade atributos, filtra, agrupa, muestrea) y exporta a uno o varios destinos. Ventajas: los servicios solo conocen un endpoint; cambiar de Jaeger a Tempo es cambiar el Collector; el muestreo de cola se hace ahí.
# km0/observabilidad/otel-collector.yaml
receivers:
otlp:
protocols:
grpc:
endpoint: 0.0.0.0:4317
tls:
cert_file: /certs/otel-collector.crt # emitidos por Vault PKI (06-04)
key_file: /certs/otel-collector.key
client_ca_file: /certs/ca.crt
processors:
batch:
timeout: 5s
memory_limiter:
limit_mib: 512
check_interval: 1s
# Muestreo de cola: decide con la traza completa. Conserva todas las que tengan error
# o duren más de 500 ms (el SLO de 07-01), y un 10 % del resto.
tail_sampling:
decision_wait: 10s
policies:
- name: errores
type: status_code
status_code: {status_codes: [ERROR]}
- name: lentas
type: latency
latency: {threshold_ms: 500}
- name: resto
type: probabilistic
probabilistic: {sampling_percentage: 10}
exporters:
otlp/tempo:
endpoint: tempo:4317
tls: {insecure: true} # red interna del compose; en producción, mTLS
# otlp/jaeger: # alternativa equivalente
# endpoint: jaeger:4317
service:
pipelines:
traces:
receivers: [otlp]
processors: [memory_limiter, tail_sampling, batch]
exporters: [otlp/tempo]# km0/docker-compose.yml (fragmento)
otel-collector:
image: otel/opentelemetry-collector-contrib:0.108.0
command: ["--config=/etc/otelcol/config.yaml"]
volumes:
- ./observabilidad/otel-collector.yaml:/etc/otelcol/config.yaml:ro
- ./certs/otel-collector:/certs:ro
ports: ["4317:4317"]
tempo:
image: grafana/tempo:2.5.0
command: ["-config.file=/etc/tempo/tempo.yaml"]
volumes:
- ./observabilidad/tempo.yaml:/etc/tempo/tempo.yaml:ro
# tempo.yaml apunta el almacenamiento de bloques al bucket km0-trazas de MinIOBackends
| Backend | Modelo | Almacenamiento | Integración | Cuándo |
|---|---|---|---|---|
| Jaeger | El clásico de CNCF; UI propia muy madura para explorar trazas | Cassandra, Elasticsearch, o su almacén propio | Recibe OTLP; datasource en Grafana | Equipos que quieren la UI de Jaeger o ya tienen Cassandra/ES |
| Tempo | Solo almacena e indexa por trace_id (y TraceQL para buscar por atributos) |
Almacenamiento de objetos (MinIO, S3): muy barato | Nativo en Grafana; enlaces desde Loki y Prometheus | Pila Grafana, mucho volumen, presupuesto contenido |
| Zipkin | El pionero (Twitter); formato propio que OTel también exporta | Cassandra, ES, MySQL | UI propia | Sistemas heredados que ya lo usan |
Kilómetro Cero usa Tempo, con los bloques en el bucket km0-trazas de MinIO, por coherencia con Loki y Grafana.
Muestreo: head frente a tail
Guardar todas las trazas de 100 peticiones/s con 8-12 spans cada una es caro y casi inútil: la mayoría son idénticas y correctas. Dos estrategias:
| Muestreo | Dónde decide | Con qué información | Ventaja | Inconveniente |
|---|---|---|---|---|
Head (TraceIdRatioBased en el SDK) |
En el primer servicio, al crear la traza | Ninguna: es una probabilidad | Barato; los servicios que no muestrean ni siquiera crean spans | Descarta el 90 % de los errores, que es justo lo que interesa |
Tail (tail_sampling en el Collector) |
En el Collector, cuando la traza está completa | Duración, estado, atributos | Conserva el 100 % de las trazas con error o lentas | Los servicios envían todo; el Collector necesita memoria para esperar 10 s por traza |
Lo habitual, y lo que hace km0/, es combinarlos: un ParentBased en el SDK para que la decisión sea coherente en toda la traza, un ratio alto en origen (o 100 % en servicios de bajo tráfico) y tail sampling en el Collector para quedarse con lo interesante.
- Una traza de "crear pedido" comentada
Así aparece en Grafana (Tempo) la traza 4bf92f35... del pedido P-2026-000125 de Ana, que tardó 482 ms y acabó bien. La sangría indica la relación padre-hijo:
| Span | Servicio | Inicio (ms) | Duración (ms) | Atributos relevantes / observaciones |
|---|---|---|---|---|
kong.proxy POST /api/v1/pedidos |
Kong | 0 | 482 | http.status_code=201, km0.request_id=c1d2... |
⎿ POST /pedidos |
pedidos-2 | 3 | 476 | km0.pedido_id=P-2026-000125, enduser.pseudo=u:9f1c3a |
⎿ saga.reservar_stock |
pedidos-2 | 6 | 118 | span manual de saga_pedido.py |
⎿ km0.Inventario/ReservarStock (cliente) |
pedidos-2 | 7 | 116 | rpc.grpc.status_code=0, deadline 300 ms |
⎿ km0.Inventario/ReservarStock (servidor) |
inventario-1 | 9 | 111 | net.peer.name=pedidos-2 (mTLS, SPIFFE) |
⎿ SELECT km0_inventario |
inventario-1 | 11 | 4 | db.statement=SELECT ... FROM stock WHERE producto=$1 FOR UPDATE |
⎿ UPDATE km0_inventario |
inventario-1 | 16 | 98 | db.statement=UPDATE stock SET ...: espera de bloqueo |
⎿ COMMIT |
inventario-1 | 115 | 5 | |
⎿ saga.cobrar |
pedidos-2 | 126 | 331 | |
⎿ km0.Pagos/Cobrar (cliente → servidor) |
pedidos-2 → pagos | 127 | 329 | |
⎿ POST pasarela.example/charges |
pagos | 131 | 318 | http.status_code=200: la pasarela externa domina la latencia |
⎿ INSERT km0_pedidos (Cassandra) |
pedidos-2 | 459 | 9 | db.system=cassandra, consistencia LOCAL_QUORUM |
⎿ publicar pedidos.eventos |
pedidos-2 | 469 | 8 | messaging.destination=pedidos.eventos |
⎿ consumir pedidos.eventos |
analitica | 612 | 21 | hijo de "publicar": empieza 140 ms después de que Ana ya tenga su 201 |
Lo que la traza cuenta y las métricas no podían:
- De los 482 ms, 318 son la pasarela de pagos externa. Nada en
km0/puede acelerarlos; sí se puede paralelizar la reserva y el cobro, o aceptar el pedido "pendiente de confirmar" (07-04). - El
UPDATEdeinventariotardó 98 ms en una operación que normalmente tarda 2: hay contención de bloqueo sobrevino-crianza(la Semana de la Vendimia, muchos pedidos del mismo producto). La métrica p99 deReservarStocklo mostraba subiendo; la traza dice en qué sentencia. - El span de
analiticapertenece a la misma traza aunque ocurra después de la respuesta: es el efecto de propagartraceparenten las cabeceras Kafka. - Cada span lleva el
trace_id, y por tanto un clic lleva a las líneas de log de los cinco servicios con ese identificador, en orden.
- Unir las tres señales
La observabilidad no son tres herramientas, sino un recorrido: la métrica dice que algo va mal, la traza dice dónde, el log dice por qué. Grafana permite que ese recorrido sean tres clics:
- Exemplars: Prometheus puede guardar, junto a cada bucket del histograma, el
trace_idde una observación reciente. Enprometheus_clientse pasaexemplar={"trace_id": ...}aobserve(); en el gráfico del p99, cada punto tiene un rombo que abre esa traza. Es el enlace métrica → traza. - Enlace log ↔ traza: el datasource de Loki en Grafana se configura con una derived field que reconoce
trace_iden el JSON y muestra un botón "ver traza en Tempo"; y Tempo, a la inversa, con "logs de esta traza" que lanza{servicio=~".+"} | json | trace_id="..."en Loki. - Traza → métrica: Tempo puede generar métricas RED a partir de los spans (span metrics), útil para servicios sin instrumentar con Prometheus.
# km0/observabilidad/grafana/provisioning/datasources/observabilidad.yml
apiVersion: 1
datasources:
- name: Loki
type: loki
url: http://loki:3100
jsonData:
derivedFields:
- name: trace_id
matcherRegex: '"trace_id":\s*"(\w+)"'
url: "$${__value.raw}"
datasourceUid: tempo
- name: Tempo
type: tempo
uid: tempo
url: http://tempo:3200
jsonData:
tracesToLogsV2:
datasourceUid: loki
filterByTraceID: true
tags: [{key: "service.name", value: "servicio"}]
tracesToMetrics:
datasourceUid: prometheusUn ejemplo de exemplar en metricas.py (07-01), ahora que existe el contexto de traza:
# km0/servicios/comun/metricas.py (añadido)
def _exemplar_actual() -> dict | None:
ctx = trace.get_current_span().get_span_context()
return {"trace_id": format(ctx.trace_id, "032x")} if ctx.is_valid and ctx.trace_flags.sampled else None
# dentro de medir(), en el finally:
PETICION_DURACION.labels(...).observe(duracion, exemplar=_exemplar_actual())Con esto, el sábado a las 10:13 el recorrido de Jordi es: alerta PedidosBurnRateRapido → panel con la tasa de error → rombo de exemplar sobre el pico → traza de un pedido fallido con el span ReservarStock en rojo y rpc.grpc.status_code=UNAVAILABLE → "logs de esta traza" → {"evento": "circuito_abierto", "dependencia": "inventario"}... que es donde entra el resto del módulo.
Errores Comunes y Consejos
- Logs de texto libre "porque se leen mejor". Se leen mejor uno a uno; no se consultan. Nombre de evento estable + campos; para leer en local,
structlogtiene un renderizador de consola con colores. - Promocionar
trace_idopedido_ida etiqueta de Loki. Cada valor crea un stream y Loki se degrada igual que Prometheus con cardinalidad alta. Etiquetas:servicio,instancia,nivel,entorno. Lo demás, en el JSON. - Loguear cabeceras o cuerpos completos en
debugy dejarlo activado. Es la fuga de tokens y datos personales más habitual. Un processor de censura como red, ydebugapagado por defecto. - Perder el contexto en los saltos asíncronos. Un
ThreadPoolExecutor, unasyncio.create_tasko un mensaje Kafka sinpropagate.injectrompen la traza. Ante una traza que "termina de repente", buscar el salto sin propagación. - Muestreo solo en cabecera. Descarta el 90 % de los errores. Tail sampling en el Collector para quedarse con errores y trazas lentas.
- Trazas sin atributos de negocio. Un span
POST /pedidossinkm0.pedido_idobliga a buscar por tiempo. Añadir el identificador de dominio como atributo (en trazas y logs no hay problema de cardinalidad). - Confundir logs operativos con auditoría. Retención, inmutabilidad y acceso son distintos (06-05). Ni la auditoría va a Loki con 14 días, ni los logs de depuración van al bucket con object lock.
- Enviar trazas directamente al backend desde cada servicio. Acopla todos los servicios al backend y hace imposible el tail sampling. Un Collector en medio.
Ejercicios
Ejercicio 1. Marc reporta a soporte que su pedido de la Semana del Queso Artesano "dio error" y aporta el X-Request-Id que la web le mostró: 7a1b.... Marta abre Grafana. (a) Escribe la consulta LogQL que reúne todo lo que pasó con esa petición en todos los servicios. (b) La consulta devuelve líneas de Kong y pedidos con trace_id 9c4d..., pero ninguna de inventario; en cambio, hay líneas de inventario con un trace_id distinto en el mismo segundo. ¿Qué se ha roto y dónde lo buscarías? (c) Marta encuentra en la línea de pedidos el campo "usuario": "u:2b77e1". ¿Cómo confirma que es Marc sin que nadie más pueda hacerlo por su cuenta, y por qué el log no contiene su email?
Ejercicio 2. El equipo decide que reparto publique una línea de log por cada posición recibida (140 repartidores, una cada 5 s) con el nivel info, y que debug incluya el cuerpo del mensaje. (a) Calcula cuántas líneas al día genera solo reparto y estima el volumen si cada línea ocupa 350 bytes. (b) Propón una alternativa que conserve la capacidad de investigar "por qué el panel mostraba a furgoneta-3-017 en Girona cuando estaba en Lleida" sin ese volumen. (c) ¿Qué problema de privacidad añade el debug con el cuerpo, y con qué mecanismo de esta lección lo mitigarías como último recurso?
Ejercicio 3. La traza de la sección 8 muestra 318 ms en la pasarela de pagos y 98 ms de espera de bloqueo en inventario. Un compañero propone "poner TraceIdRatioBased(1.0) en todos los servicios para no perder ninguna traza como esta". (a) Calcula cuántos spans por segundo recibiría el Collector con 100 peticiones/s a pedidos y 12 spans por traza, más el tráfico de catalogo (400 peticiones/s, 4 spans). (b) Explica por qué la configuración de otel-collector.yaml de la sección 7 ya habría conservado esta traza aunque el ratio fuera 0,1, y qué habría pasado si el ratio 0,1 se hubiera configurado sin ParentBased. (c) Un span de la compensación de la saga, ejecutado dos horas después de que el pedido fallara, aparece como traza nueva y no dentro de la del pedido. ¿Qué falta y dónde?
Soluciones
Ejercicio 1.
(a) {servicio=~".+"} | json | request_id="7a1b...", con el rango de tiempo acotado al momento del pedido (Loki escanea el contenido, así que el rango importa). (b) Kong y pedidos comparten trace_id, luego la propagación HTTP funciona; inventario tiene otro trace_id en el mismo instante, luego inventario está creando trazas nuevas en lugar de continuar la de pedidos: el traceparent no llega en los metadatos gRPC o no se lee. Sospechosos por orden: GrpcInstrumentorClient no instrumentado en pedidos (el cliente no inyecta), un interceptor propio en inventario que reconstruye los metadatos y descarta traceparent (revisar AuthInterceptor/ServicioInterceptor de 06-04), o un stub gRPC creado antes de llamar a configurar_trazas. Se comprueba con un log.debug de los metadatos recibidos en inventario (censurando authorization) o mirando en Tempo si el span servidor de inventario es raíz. (c) La seudonimización de 06-02 es una función con clave (HMAC) sobre el identificador del usuario: Marta pide al servicio de identidades, con su rol de operador y su ticket de soporte (queda auditado, 06-05), que resuelva u:2b77e1; o calcula el seudónimo del sub de Marc y compara. El log no contiene el email porque es un dato personal que viaja a Loki, a MinIO y a la pantalla de cualquiera con acceso a Grafana: se registra el seudónimo, que no identifica sin la clave, y el pedido_id, que es lo que soporte necesita.
Ejercicio 2.
(a) 140 × (86 400 / 5) = 2 419 200 líneas al día; a 350 bytes, unos 847 MB diarios sin comprimir solo de posiciones, unos 25 GB al mes. (b) Loguear en info solo lo anómalo (posiciones rechazadas por validación, saltos imposibles entre dos posiciones, mensajes fuera de orden) y un resumen por repartidor y minuto (posiciones_recibidas, ultima_ts); para reconstruir la ruta de furgoneta-3-017, los datos ya están en reparto.posiciones con retención de horas y en el lago HDFS (04-02): investigar consultando esos datos, no el log. Además, una traza muestreada por repartidor cada N mensajes da el recorrido interno sin loguear cada uno. (c) El cuerpo contiene la posición exacta de un empleado con su identificador, un dato personal (06-05, ejercicio 3); si alguien activa debug en producción, se copia a Loki sin las garantías de acceso de reparto.posiciones. Como último recurso, el processor censurar con posicion/lat/lon en CLAVES_PROHIBIDAS (y que el cuerpo se loguee como campos, nunca como cadena opaca, para que el censor pueda actuar); pero la política correcta es no loguear el cuerpo.
Ejercicio 3.
(a) pedidos: 100 × 12 = 1 200 spans/s; catalogo: 400 × 4 = 1 600 spans/s; total unos 2 800 spans/s, 242 millones al día, que el Collector debe recibir y retener 10 s cada uno para decidir (unos 28 000 spans en memoria de forma permanente, factible), y de los que Tempo escribiría lo que las políticas dejen pasar. Con TraceIdRatioBased(1.0) el coste está en la red y en el Collector, no en el almacenamiento, siempre que haya tail sampling. (b) La traza dura 482 ms y la política lentas conserva todo lo que supere 500 ms... no la conservaría por duración; pero tampoco importa: si tuvo un span con error, la política errores la guarda; si no, cae en el 10 % probabilístico. Con el ratio 0,1 en cabecera, la decisión la tomó Kong (o pedidos) al crear la traza, y ParentBased hace que todos los servicios la respeten: o se muestrea entera o ninguno la genera. Sin ParentBased, cada servicio tiraría su propio dado: pedidos la conservaría y inventario no, y la traza quedaría con agujeros (el span ReservarStock servidor ausente, justo el de los 98 ms). (c) Falta guardar el traceparent del pedido en la fila de la tabla sagas cuando se crea la saga, y al ejecutar la compensación reconstruir el contexto con propagate.extract sobre ese valor y crear el span de compensación con context=ctx (o como link a la traza original si se prefiere una traza nueva enlazada). Está en saga_pedido.py, en el punto donde el relay o el planificador retoman una saga pendiente.
Conclusión
Con esta lección la observabilidad de Kilómetro Cero tiene sus tres señales. Los logs dejan de ser texto en ficheros de contenedores efímeros para ser eventos JSON con campos comunes (ts en UTC, nivel, servicio, instancia, trace_id, request_id, usuario seudonimizado, evento), sin secretos ni datos personales, recogidos por Promtail en cada nodo y almacenados en Loki con etiquetas de baja cardinalidad, consultables con LogQL por el X-Request-Id que Kong asigna o por el trace_id que OpenTelemetry inyecta. Las trazas dan el recorrido: spans anidados con traceparent W3C propagado por HTTP y gRPC sin código, y a mano por las cabeceras de Kafka y por la tabla sagas, enviados a un Collector que muestrea en cola para conservar errores y lentitud, y almacenados en Tempo. Y las tres se unen en Grafana con exemplars y enlaces log↔traza, de modo que una alerta se convierte en tres clics en una traza con el span culpable y las líneas de log que lo explican.
Ya sabemos que falla, dónde y por qué. El recorrido de Jordi terminaba en un log de pedidos que decía circuito_abierto hacia inventario, y en Tempo un span ReservarStock en rojo con UNAVAILABLE: inventario-1 no responde. La pregunta siguiente no es de observabilidad, sino de supervivencia: ¿cómo detecta la plataforma que un nodo ha caído y no simplemente está lento? ¿Quién decide que inv-vlc pase a ser primario, y cómo se evita que inv-bcn vuelva creyendo que sigue siéndolo? ¿Y si lo que se ha perdido no es un proceso sino los datos de km0_inventario de las últimas dos horas? La siguiente lección trata la gestión de fallos y la recuperación: detección, failover con cercado, backups y PITR, y qué hacer en los minutos en que todo lo anterior está pasando.
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
