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

  1. Por qué tail -f ya no vale
  2. Logs estructurados: campos estándar, niveles y qué no loguear
  3. Pipeline de recogida: agente, almacén y consulta
  4. servicios/comun/logs.py: JSON con trace_id inyectado
  5. Trazabilidad distribuida: trazas, spans y contexto
  6. Instrumentación con OpenTelemetry: servicios/comun/trazas.py
  7. Collector, backends y muestreo
  8. Una traza de "crear pedido" comentada
  9. Unir las tres señales
  10. Errores comunes y consejos
  11. Ejercicios y soluciones
  12. Conclusión

  1. Por qué tail -f ya no vale

En 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-000125 genera líneas en Kong, pedidos, inventario, pagos y analitica. 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 de tomate-rosa en 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.

  1. 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 Authorization completas. Una traza HTTP en debug que 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) y pedido_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.

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

  1. servicios/comun/logs.py: JSON con trace_id inyectado

Con 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:

  • structlog no formatea cadenas: log.warning("reserva_stock_rechazada", producto=...) produce claves separadas. Escribir log.warning(f"Rechazado {producto}") funciona, pero pierde la estructura. La disciplina es "nombre de evento estable + campos".
  • El processor censurar es una red de seguridad, no la política: la política es no pasar esos datos. Pero cuesta poco y evita que un log.debug("peticion", headers=dict(request.headers)) de un martes por la noche acabe volcando tokens.
  • El trace_id se 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 mismo trace_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.

  1. 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_id de 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. Tiene span_id y un parent_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

  1. Instrumentación con OpenTelemetry: servicios/comun/trazas.py

OpenTelemetry (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 pedidos

Dos aclaraciones:

  • gRPC, HTTP y PostgreSQL no requieren código: los instrumentors interceptan las llamadas, crean spans y meten/leen traceparent en los metadatos. En inventario, AuthInterceptor y el MetricasInterceptor de 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.py guarda traceparent en la fila de la tabla sagas, y cuando el relay del outbox o una compensación retoman el trabajo, reconstruyen el contexto con propagate.extract sobre ese valor. Sin esto, la compensación de P-2026-000125 aparecería como una traza huérfana sin relación con el pedido.

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

Backends

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.

  1. 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 UPDATE de inventario tardó 98 ms en una operación que normalmente tarda 2: hay contención de bloqueo sobre vino-crianza (la Semana de la Vendimia, muchos pedidos del mismo producto). La métrica p99 de ReservarStock lo mostraba subiendo; la traza dice en qué sentencia.
  • El span de analitica pertenece a la misma traza aunque ocurra después de la respuesta: es el efecto de propagar traceparent en 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.

  1. 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_id de una observación reciente. En prometheus_client se pasa exemplar={"trace_id": ...} a observe(); 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_id en 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: prometheus

Un 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, structlog tiene un renderizador de consola con colores.
  • Promocionar trace_id o pedido_id a 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 debug y dejarlo activado. Es la fuga de tokens y datos personales más habitual. Un processor de censura como red, y debug apagado por defecto.
  • Perder el contexto en los saltos asíncronos. Un ThreadPoolExecutor, un asyncio.create_task o un mensaje Kafka sin propagate.inject rompen 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 /pedidos sin km0.pedido_id obliga 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

Módulo 2: Comunicación en Sistemas Distribuidos

Módulo 3: Consistencia y Replicación

Módulo 4: Almacenamiento Distribuido

Módulo 5: Computación Distribuida

Módulo 6: Seguridad en Sistemas Distribuidos

Módulo 7: Monitoreo y Mantenimiento

Módulo 8: Casos de Estudio y Aplicaciones

© Copyright 2026. Todos los derechos reservados