La lliçó anterior va deixar en Jordi Sala amb una alerta a la mà: ComandesBurnRateRapid, el 8 % de les comandes falla. Les mètriques li diuen quant i des de quan; no li diuen quines comandes, ni a quin servei, ni per què. Al monòlit de Quilòmetre Zero la resposta era a un grep de distància en un únic fitxer de log. Avui una comanda travessa Kong, tres rèpliques de comandes, dues d'inventari, pagaments, un tòpic de Kafka i analitica, cadascun escrivint al seu propi contenidor; el log de la comanda P-2026-000125 està repartit en set llocs i en cap no es diu igual. Aquesta lliçó construeix els altres dos pilars de l'observabilitat: logs estructurats i centralitzats, perquè un identificador reuneixi tot el que va passar amb una petició, i traces distribuïdes, per veure el recorregut i el temps d'aquesta petició per tots els serveis. Al final, els tres senyals queden units a Grafana: de la mètrica que alerta, a la traça que localitza, al log que explica.
Contingut
- Per què
tail -fja no serveix - Logs estructurats: camps estàndard, nivells i què no cal registrar
- Pipeline de recollida: agent, magatzem i consulta
serveis/comu/logs.py: JSON ambtrace_idinjectat- Traçabilitat distribuïda: traces, spans i context
- Instrumentació amb OpenTelemetry:
serveis/comu/traces.py - Collector, backends i mostreig
- Una traça de "crear comanda" comentada
- Unir els tres senyals
- Errors comuns i consells
- Exercicis i solucions
- Conclusió
- Per què
tail -f ja no serveix
tail -f ja no serveixEn una màquina, un log és un fitxer i un incident s'investiga amb tail -f, grep i less. En un sistema distribuït això deixa de funcionar per raons que no són de comoditat, sinó estructurals:
- La informació d'una petició està repartida. Crear la comanda
P-2026-000125genera línies a Kong,comandes,inventari,pagamentsianalitica. Per reconstruir la història cal obrir cinc terminals i endevinar quina línia de cadascuna correspon a aquesta comanda. - Els rellotges no coincideixen. Com es va veure a 01-05, ordenar per timestamp línies de màquines diferents pot invertir causa i efecte. Cal un identificador de correlació, no només l'hora.
- Els contenidors són efímers. Quan Kubernetes (07-05) reemplaça un pod de
comandes, el seu log desapareix amb ell. Si la fallada va passar just abans de morir, s'ha perdut l'evidència. - El format lliure no es consulta.
Comanda 125 rebutjada per estoc (tomaquet-rosa)és llegible per a una persona; per preguntar "quants rebutjos per estoc detomaquet-rosaa Lleida en la darrera hora" cal parsejar amb expressions regulars fràgils. - Volum. 140 repartidors enviant una posició cada 5 s, 100 peticions/s a
comandes: són milions de línies al dia. Ningú no les llegeix; cal poder filtrar-les.
La solució té dues parts: que cada servei escrigui logs estructurats amb camps comuns, i que un pipeline els reculli de tots els nodes en un magatzem central consultable.
- Logs estructurats: camps estàndard, nivells i què no cal registrar
Un log estructurat és un registre amb camps amb nom, normalment una línia JSON per esdeveniment:
{"ts": "2026-09-12T10:13:41.208Z", "nivell": "warning", "servei": "comandes", "instancia": "comandes-2",
"trace_id": "4bf92f3577b34da6a3ce929d0e0e4736", "span_id": "00f067aa0ba902b7",
"request_id": "c1d2e3f4-5a6b-7c8d-9e0f-1a2b3c4d5e6f", "usuari": "u:9f1c3a",
"esdeveniment": "reserva_estoc_rebutjada", "comanda_id": "P-2026-000125",
"producte": "vi-crianca", "sollicitades": 6, "disponibles": 2, "replica": "inv-bcn"}Els camps que tots els serveis de km0/ inclouen sempre:
| Camp | Contingut | Per què |
|---|---|---|
ts |
Instant en UTC, ISO 8601 amb mil·lisegons | Sense zona horària no hi ha ordenació possible entre nodes; UTC evita el canvi d'hora |
nivell |
debug, info, warning, error, critical |
Filtrar i alertar |
servei |
comandes, inventari, ... |
Primer filtre a qualsevol consulta |
instancia |
comandes-2, nom del pod |
Distingir una rèplica malalta |
trace_id |
Identificador de la traça OpenTelemetry (32 hex) | Correlació entre serveis i amb les traces |
span_id |
Span actual (16 hex) | Enllaçar la línia amb el pas concret de la traça |
request_id |
L'X-Request-Id que Kong genera i propaga (06-05) |
És l'identificador que veu el client i que l'Anna pot donar a suport |
usuari |
Identificador pseudonimitzat del subjecte (06-02) | Investigar sense exposar la identitat |
esdeveniment |
Nom curt i estable del succés (comanda_creada, reserva_estoc_rebutjada) |
Comptar i filtrar per tipus sense parsejar el missatge |
missatge (opcional) |
Text per a humans | Context addicional |
Els camps específics (comanda_id, producte, replica) s'afegeixen com a claus addicionals, no dins del missatge. I sí: aquí comanda_id és benvingut. El que a les mètriques era un problema de cardinalitat (07-01) als logs és exactament el que es busca, perquè un log s'indexa per text o per unes poques etiquetes, no crea una sèrie per valor.
Nivells
| Nivell | Quan | Exemple a comandes |
|---|---|---|
debug |
Detall per a desenvolupament; apagat en producció llevat d'investigació puntual | Cos de la petició gRPC a inventari |
info |
Fites normals del negoci | comanda_creada, saga_completada |
warning |
Alguna cosa inesperada que el sistema ha gestionat | reserva_estoc_rebutjada, reintent d'una crida |
error |
Una operació ha fallat i algú hauria de mirar-s'ho | pagament_error_tecnic, compensacio_fallida |
critical |
El servei no pot continuar | No connecta amb km0_comandes en arrencar |
Què no cal registrar
Els logs viatgen per la xarxa, s'emmagatzemen mesos i els llegeix molta gent. Tot el que el Mòdul 6 va protegir amb xifratge i control d'accés pot acabar en clar en un log si no es vigila:
- Secrets: tokens JWT (ni tan sols "per depurar"), contrasenyes, credencials de Vault, claus d'API de la passarel·la de pagaments, capçaleres
Authorizationcompletes. Una traça HTTP endebugque aboqui capçaleres és una fuita. - Dades personals en clar: nom, email, telèfon i adreça de l'Anna. Es registra
usuari: "u:9f1c3a"(el pseudònim de 06-02) icomanda_id; qui necessiti el telèfon l'obté de la base de dades amb el seu permís i la seva auditoria. - Dades de targeta: mai, de cap manera (06-02).
- Cossos complets de peticions: a més de dades personals, són enormes. Es registren els camps rellevants triats un a un.
Els logs d'auditoria de 06-05 (qui va fer què, hash encadenat, object lock a MinIO) són un canal diferent amb garanties diferents: immutabilitat i retenció llarga. Els logs operatius d'aquesta lliçó són per diagnosticar, es retenen setmanes i es poden esborrar; no s'han de barrejar.
- Pipeline de recollida: agent, magatzem i consulta
flowchart LR
subgraph node1["Node 1"]
P1[comandes-1<br/>stdout JSON]
I1[inventari-1<br/>stdout JSON]
A1[Promtail]
P1 --> A1
I1 --> A1
end
subgraph node2["Node 2"]
P2[comandes-2]
K[Kong]
A2[Promtail]
P2 --> A2
K --> A2
end
L[(Loki<br/>index per etiquetes<br/>chunks a MinIO km0-logs)]
G[Grafana<br/>LogQL]
A1 --> L
A2 --> L
G --> L
Cada servei escriu JSON a stdout (no en fitxers propis: el contenidor no ha de saber on acaben els seus logs). Un agent per node llegeix la sortida de tots els contenidors del node, hi afegeix etiquetes (servei, instància, node) i envia al magatzem central. La consulta es fa des d'una interfície que entén el format.
| Component | Opció "Loki" | Opció "ELK" | Altres |
|---|---|---|---|
| Agent | Promtail (o Grafana Alloy) | Filebeat, Logstash | Fluent Bit / Fluentd (totes dues piles), Vector |
| Magatzem | Loki | Elasticsearch / OpenSearch | — |
| Consulta | Grafana (LogQL) | Kibana / OpenSearch Dashboards (KQL, Lucene) | — |
Loki contra Elasticsearch, que és la decisió d'arquitectura real:
| Aspecte | Loki | Elasticsearch |
|---|---|---|
| Què indexa | Només les etiquetes (servei, instancia, nivell); el contingut es desa comprimit sense indexar |
El text complet de cada camp (índex invertit) |
| Cost d'emmagatzematge | Baix (chunks comprimits en un bucket d'objectes, p. ex. MinIO 04-03) | Alt (l'índex pot ocupar més que les dades) |
| Cost d'ingesta | Baix | Alt en CPU i memòria |
Consulta "tot el de comandes amb trace_id=X en la darrera hora" |
Filtra per etiqueta i després escaneja el contingut d'aquella hora: ràpid si el rang és acotat | Instantani gràcies a l'índex |
| Cerca lliure en mesos de dades | Lenta (escaneig) | Ràpida |
| Analítica sobre logs (agregacions complexes) | Limitada (LogQL té mètriques sobre logs, però no és un motor analític) | Potent |
| Integració | Nativa amb Grafana, Prometheus i Tempo: mateix llenguatge d'etiquetes, enllaços log↔traça↔mètrica | Kibana; APM propi |
| Quan | Diagnòstic operatiu amb etiquetes conegudes, pressupost contingut, pila Grafana | Cerca exploratòria intensa, compliment normatiu amb cerques ad hoc, equips que ja operen Elasticsearch |
Quilòmetre Zero tria Loki perquè el seu patró de consulta és gairebé sempre "servei + rang + trace_id", l'emmagatzematge va a MinIO i ja fa servir Grafana i Prometheus. La cardinalitat torna a importar, aquest cop a les etiquetes de Loki: servei, instancia i nivell són etiquetes; trace_id, comanda_id i usuari no ho són (crearien un stream per valor), es busquen dins del contingut amb | json | trace_id="...".
Retenció i cost
Els logs creixen sense límit si ningú no decideix el contrari. Una política típica: debug no s'envia en producció; info es conserva 14 dies; warning/error 90 dies; els logs de Kong (una línia per petició) 30 dies; els d'auditoria són un altre sistema amb anys de retenció. A Loki la retenció es configura per stream i els chunks antics s'esborren del bucket. Reduir el volum en origen (no registrar cada posició de repartiment, sinó un resum per minut i els rebutjos) sol ser la mesura més eficaç.
Configuració mínima de Promtail i Loki a docker-compose.yml:
# km0/docker-compose.yml (fragment)
loki:
image: grafana/loki:3.1.0
command: ["-config.file=/etc/loki/loki.yaml"]
volumes:
- ./observabilitat/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:
- ./observabilitat/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/observabilitat/promtail.yaml
server:
http_listen_port: 9080
clients:
- url: http://loki:3100/loki/api/v1/push
scrape_configs:
- job_name: contenidors
docker_sd_configs:
- host: unix:///var/run/docker.sock
relabel_configs:
# L'etiqueta de compose "com.docker.compose.service" passa a ser l'etiqueta "servei"
- source_labels: [__meta_docker_container_label_com_docker_compose_service]
target_label: servei
- source_labels: [__meta_docker_container_name]
regex: "/(.*)"
target_label: instancia
pipeline_stages:
- json:
expressions:
nivell: nivell
- labels:
nivell: # només "nivell" es promociona a etiqueta; trace_id es queda al contingut# km0/observabilitat/loki.yaml (fragment rellevant)
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 dies per defecte
compactor:
retention_enabled: true
serveis/comu/logs.py: JSON amb trace_id injectat
serveis/comu/logs.py: JSON amb trace_id injectatAmb structlog, cada servei configura un sol cop el logger i després fa servir log.info("esdeveniment", clau=valor). Un processor consulta el context d'OpenTelemetry (secció 6) i afegeix trace_id i span_id si hi ha un span actiu:
# km0/serveis/comu/logs.py
"""Logs JSON estructurats amb correlació per 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 es desa en una ContextVar: és local a la petició
# en curs (funciona amb fils i amb asyncio) i qualsevol log el pot llegir.
request_id_actual: ContextVar[str | None] = ContextVar("request_id", default=None)
SERVEI = os.environ.get("KM0_SERVEI", "desconegut")
INSTANCIA = os.environ.get("HOSTNAME", "local")
# Claus que mai no han de sortir en un log, encara que algú les passi per descuit.
CLAUS_PROHIBIDES = {"authorization", "password", "contrasenya", "token", "targeta", "cvv", "telefon", "email"}
def injectar_context(logger, metode, esdeveniment: dict) -> dict:
"""Processor de structlog: afegeix servei, instancia, trace_id, span_id i request_id."""
esdeveniment["servei"] = SERVEI
esdeveniment["instancia"] = INSTANCIA
span = trace.get_current_span()
ctx = span.get_span_context()
if ctx.is_valid:
# format(x, "032x") = 32 dígits hexadecimals, el format W3C
esdeveniment["trace_id"] = format(ctx.trace_id, "032x")
esdeveniment["span_id"] = format(ctx.span_id, "016x")
rid = request_id_actual.get()
if rid:
esdeveniment["request_id"] = rid
return esdeveniment
def censurar(logger, metode, esdeveniment: dict) -> dict:
"""Processor: substitueix el valor de qualsevol clau sensible per '[censurat]'."""
for clau in list(esdeveniment):
if clau.lower() in CLAUS_PROHIBIDES:
esdeveniment[clau] = "[censurat]"
return esdeveniment
def configurar_logs(nivell: str = "INFO") -> None:
structlog.configure(
processors=[
structlog.contextvars.merge_contextvars, # camps lligats amb bind_contextvars
structlog.processors.add_log_level, # -> "level"
structlog.processors.TimeStamper(fmt="iso", utc=True, key="ts"),
injectar_context,
censurar,
structlog.processors.EventRenamer("esdeveniment"), # el primer argument passa a "esdeveniment"
structlog.processors.JSONRenderer(),
],
wrapper_class=structlog.make_filtering_bound_logger(getattr(logging, nivell)),
logger_factory=structlog.PrintLoggerFactory(file=sys.stdout),
)
log = structlog.get_logger()Com es fa servir a comandes, amb el middleware que captura l'X-Request-Id que Kong propaga:
# km0/serveis/comandes/app.py (fragment)
from serveis.comu.logs import log, request_id_actual, configurar_logs
configurar_logs(nivell=os.environ.get("KM0_LOG_NIVELL", "INFO"))
@app.middleware("http")
async def middleware_request_id(request: Request, call_next):
rid = request.headers.get("x-request-id", "sense-request-id")
token = request_id_actual.set(rid)
try:
return await call_next(request)
finally:
request_id_actual.reset(token)
# A la saga (03-05), quan es rebutja una reserva:
log.warning("reserva_estoc_rebutjada", comanda_id=comanda.id, producte=linia.producte,
sollicitades=linia.quantitat, disponibles=resp.disponibles, replica=resp.replica)Detalls que convé notar:
structlogno formata cadenes:log.warning("reserva_estoc_rebutjada", producte=...)produeix claus separades. Escriurelog.warning(f"Rebutjat {producte}")funciona, però perd l'estructura. La disciplina és "nom d'esdeveniment estable + camps".- El processor
censurarés una xarxa de seguretat, no la política: la política és no passar aquestes dades. Però costa poc i evita que unlog.debug("peticio", headers=dict(request.headers))d'un dimarts a la nit acabi abocant tokens. - El
trace_ides llegeix del span actiu d'OpenTelemetry. Sense traces configurades, el context és invàlid i simplement no apareix; amb elles, cada línia de log de qualsevol servei que participi en la mateixa petició porta el mateixtrace_id, encara que els serveis no hagin fet res explícit per passar-lo.
Consultes LogQL a Grafana:
# Tot el que va passar amb una traça concreta, a tots els serveis, en ordre
{servei=~"comandes|inventari|pagaments|analitica"} | json | trace_id="4bf92f3577b34da6a3ce929d0e0e4736"
# Rebutjos d'estoc per producte en la darrera hora (mètrica derivada de logs)
sum by (producte) (count_over_time({servei="comandes"} | json | esdeveniment="reserva_estoc_rebutjada" [1h]))
# Errors de comandes-2 que no són de la passarel·la de pagaments
{servei="comandes", instancia="comandes-2", nivell="error"} | json | esdeveniment != "pagament_error_tecnic"
# El que l'Anna reporta a suport: el seu request_id
{servei=~".+"} | json | request_id="c1d2e3f4-5a6b-7c8d-9e0f-1a2b3c4d5e6f"La segona consulta mostra un ús valuós: mètriques a partir de logs per a preguntes que no mereixen una mètrica pròpia (07-01, exercici 1) però que sí que es volen graficar de tant en tant.
- Traçabilitat distribuïda: traces, spans i context
Els logs amb trace_id responen "què va passar"; les traces responen "per on va passar i quant va trigar cada pas". El model, heretat de Dapper (Google) i estandarditzat per OpenTelemetry:
- Una traça és l'arbre complet de feina causat per una petició externa. S'identifica per un
trace_idde 128 bits. - Un span és una unitat de feina amb nom, inici, fi, atributs (
comanda_id,rpc.method,db.statement), estat (OK/ERROR) i esdeveniments. Téspan_idi unparent_span_id, llevat de l'arrel. - El context de propagació és el que viatja d'un servei a un altre perquè el span del receptor pengi de l'emissor. L'estàndard és W3C Trace Context: una capçalera
traceparent:
traceparent: 00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01
^^ ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ ^^^^^^^^^^^^^^^^ ^^
versió trace_id (32 hex) span pare flags (01 = mostrejada)(I una tracestate opcional per a dades de proveïdors.) A HTTP i gRPC és una capçalera; a Kafka, una capçalera del missatge; a la saga, es desa juntament amb l'estat a la taula sagas perquè la compensació hores després continuï penjant de la traça original.
sequenceDiagram
participant Anna
participant Kong
participant Com as comandes-2
participant Inv as inventari-1
participant Pag as pagaments
participant Kafka
participant Ana2 as analitica
Anna->>Kong: POST /api/v1/comandes
Note over Kong: crea trace_id 4bf9...<br/>span arrel "kong.proxy"
Kong->>Com: POST /comandes<br/>traceparent: 00-4bf9...-a1..-01<br/>X-Request-Id: c1d2...
Note over Com: span "POST /comandes" fill d'a1
Com->>Inv: gRPC ReservarEstoc<br/>metadata traceparent: 00-4bf9...-b2..-01
Note over Inv: span "ReservarEstoc" + span "SELECT ... FOR UPDATE"
Inv-->>Com: OK
Com->>Pag: gRPC Cobrar (traceparent ...-b2..)
Pag-->>Com: OK
Com->>Kafka: comanda.confirmada<br/>header traceparent: 00-4bf9...-c3..-01
Com-->>Kong: 201
Kong-->>Anna: 201
Kafka-->>Ana2: consumeix (140 ms després)
Note over Ana2: span "consumir comandes.esdeveniments"<br/>enllaçat a c3
- Instrumentació amb OpenTelemetry:
serveis/comu/traces.py
serveis/comu/traces.pyOpenTelemetry (OTel) aporta tres coses: un SDK que crea spans i gestiona el context, instrumentacions automàtiques per a llibreries conegudes (gRPC, FastAPI, psycopg, redis, requests, kafka-python/confluent-kafka), i un protocol d'exportació (OTLP) cap a un Collector o un backend. Els serveis de km0/ comparteixen la configuració:
# km0/serveis/comu/traces.py
"""Configuració d'OpenTelemetry per a traces distribuïdes."""
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_traces(servei: str) -> trace.Tracer:
recurs = Resource.create({
SERVICE_NAME: servei, # "comandes", "inventari"...
SERVICE_INSTANCE_ID: os.environ.get("HOSTNAME", "local"),
"deployment.environment": os.environ.get("KM0_ENTORN", "produccio"),
})
# ParentBased: si la petició arriba ja amb decisió de mostreig (flags=01), es respecta;
# si és arrel, es mostreja el 10 %. Així una traça mai no queda "a mitges".
mostreig = ParentBased(root=TraceIdRatioBased(float(os.environ.get("KM0_TRACES_RATIO", "0.1"))))
proveidor = TracerProvider(resource=recurs, sampler=mostreig)
exportador = OTLPSpanExporter(endpoint=os.environ.get("OTEL_EXPORTER_OTLP_ENDPOINT", "otel-collector:4317"),
insecure=False) # mTLS amb el certificat de Vault (06-04)
proveidor.add_span_processor(BatchSpanProcessor(exportador)) # exporta per lots, en un fil a part
trace.set_tracer_provider(proveidor)
# Autoinstrumentació: apedacen les llibreries per crear spans i propagar traceparent
GrpcInstrumentorServer().instrument()
GrpcInstrumentorClient().instrument()
PsycopgInstrumentor().instrument(enable_commenter=True) # afegeix el trace_id com a comentari SQL
RedisInstrumentor().instrument()
return trace.get_tracer(servei)
# ---- Propagació manual per capçaleres Kafka -----------------------------------
# Les capçaleres Kafka són una llista de (clau, bytes); OTel necessita un Getter/Setter
# que sàpiga llegir-les i escriure-les.
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 injectar_en_capcaleres_kafka() -> list[tuple[str, bytes]]:
"""Retorna capçaleres Kafka amb el traceparent del span actiu."""
capcaleres: list[tuple[str, bytes]] = []
propagate.inject(capcaleres, setter=_KafkaSetter())
return capcaleres
def context_des_de_capcaleres_kafka(capcaleres: list[tuple[str, bytes]] | None):
"""Reconstrueix el context de traça a partir de les capçaleres d'un missatge."""
return propagate.extract(capcaleres or [], getter=_KafkaGetter())Com ho fa servir comandes en publicar l'esdeveniment després de la saga, i analitica en consumir-lo:
# km0/serveis/comandes/esdeveniments.py (fragment)
from serveis.comu.traces import injectar_en_capcaleres_kafka
def publicar_comanda_confirmada(productor, comanda):
with tracer.start_as_current_span("publicar comandes.esdeveniments",
kind=trace.SpanKind.PRODUCER,
attributes={"messaging.system": "kafka",
"messaging.destination": "comandes.esdeveniments",
"km0.comanda_id": comanda.id}):
productor.produce("comandes.esdeveniments", key=comanda.id.encode(), value=comanda.a_json().encode(),
headers=injectar_en_capcaleres_kafka()) # traceparent viatja amb el missatge# km0/serveis/analitica/consumidor_comandes.py (fragment)
from opentelemetry import context, trace
from serveis.comu.traces import context_des_de_capcaleres_kafka
for msg in consumidor:
ctx = context_des_de_capcaleres_kafka(msg.headers())
# El span del consumidor és fill del span "publicar" de comandes: mateixa traça, 140 ms després
with tracer.start_as_current_span("consumir comandes.esdeveniments", context=ctx,
kind=trace.SpanKind.CONSUMER,
attributes={"messaging.kafka.partition": msg.partition(),
"messaging.kafka.offset": msg.offset()}):
processar(msg)
log.info("esdeveniment_processat", comanda_id=msg.key().decode()) # porta el trace_id de comandesDos aclariments:
- gRPC, HTTP i PostgreSQL no requereixen codi: els instrumentors intercepten les crides, creen spans i posen/llegeixen
traceparenta les metadades. Ainventari,AuthInterceptori elMetriquesInterceptorde 07-01 conviuen amb l'interceptor d'OTel; l'ordre no afecta el context. - Kafka sí que requereix propagació manual (o l'instrumentor de
confluent-kafka, que fa el mateix per sota): el missatge és una dada inert que es consumeix més tard, en un altre procés; ningú no el "crida". El mateix passa amb la saga:saga_comanda.pydesatraceparenta la fila de la taulasagas, i quan el relay de l'outbox o una compensació reprenen la feina, reconstrueixen el context ambpropagate.extractsobre aquest valor. Sense això, la compensació deP-2026-000125apareixeria com una traça òrfena sense cap relació amb la comanda.
- Collector, backends i mostreig
Els serveis no envien les traces al backend directament, sinó a l'OpenTelemetry Collector: un procés intermedi que rep OTLP, processa (afegeix atributs, filtra, agrupa, mostreja) i exporta a un o diversos destins. Avantatges: els serveis només coneixen un endpoint; canviar de Jaeger a Tempo és canviar el Collector; el mostreig de cua es fa allà.
# km0/observabilitat/otel-collector.yaml
receivers:
otlp:
protocols:
grpc:
endpoint: 0.0.0.0:4317
tls:
cert_file: /certs/otel-collector.crt # emesos per 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
# Mostreig de cua: decideix amb la traça completa. Conserva totes les que tinguin error
# o durin més de 500 ms (l'SLO de 07-01), i un 10 % de la resta.
tail_sampling:
decision_wait: 10s
policies:
- name: errors
type: status_code
status_code: {status_codes: [ERROR]}
- name: lentes
type: latency
latency: {threshold_ms: 500}
- name: resta
type: probabilistic
probabilistic: {sampling_percentage: 10}
exporters:
otlp/tempo:
endpoint: tempo:4317
tls: {insecure: true} # xarxa interna del compose; en producció, mTLS
# otlp/jaeger: # alternativa equivalent
# endpoint: jaeger:4317
service:
pipelines:
traces:
receivers: [otlp]
processors: [memory_limiter, tail_sampling, batch]
exporters: [otlp/tempo]# km0/docker-compose.yml (fragment)
otel-collector:
image: otel/opentelemetry-collector-contrib:0.108.0
command: ["--config=/etc/otelcol/config.yaml"]
volumes:
- ./observabilitat/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:
- ./observabilitat/tempo.yaml:/etc/tempo/tempo.yaml:ro
# tempo.yaml apunta l'emmagatzematge de blocs al bucket km0-traces de MinIOBackends
| Backend | Model | Emmagatzematge | Integració | Quan |
|---|---|---|---|---|
| Jaeger | El clàssic de la CNCF; UI pròpia molt madura per explorar traces | Cassandra, Elasticsearch, o el seu magatzem propi | Rep OTLP; datasource a Grafana | Equips que volen la UI de Jaeger o que ja tenen Cassandra/ES |
| Tempo | Només emmagatzema i indexa per trace_id (i TraceQL per buscar per atributs) |
Emmagatzematge d'objectes (MinIO, S3): molt barat | Natiu a Grafana; enllaços des de Loki i Prometheus | Pila Grafana, molt volum, pressupost contingut |
| Zipkin | El pioner (Twitter); format propi que OTel també exporta | Cassandra, ES, MySQL | UI pròpia | Sistemes heretats que ja el fan servir |
Quilòmetre Zero fa servir Tempo, amb els blocs al bucket km0-traces de MinIO, per coherència amb Loki i Grafana.
Mostreig: head contra tail
Desar totes les traces de 100 peticions/s amb 8-12 spans cadascuna és car i gairebé inútil: la majoria són idèntiques i correctes. Dues estratègies:
| Mostreig | On decideix | Amb quina informació | Avantatge | Inconvenient |
|---|---|---|---|---|
Head (TraceIdRatioBased a l'SDK) |
Al primer servei, en crear la traça | Cap: és una probabilitat | Barat; els serveis que no mostregen ni tan sols creen spans | Descarta el 90 % dels errors, que és just el que interessa |
Tail (tail_sampling al Collector) |
Al Collector, quan la traça està completa | Durada, estat, atributs | Conserva el 100 % de les traces amb error o lentes | Els serveis ho envien tot; el Collector necessita memòria per esperar 10 s per traça |
L'habitual, i el que fa km0/, és combinar-los: un ParentBased a l'SDK perquè la decisió sigui coherent a tota la traça, un ratio alt en origen (o 100 % en serveis de poc trànsit) i tail sampling al Collector per quedar-se amb el que és interessant.
- Una traça de "crear comanda" comentada
Així apareix a Grafana (Tempo) la traça 4bf92f35... de la comanda P-2026-000125 de l'Anna, que va trigar 482 ms i va acabar bé. El sagnat indica la relació pare-fill:
| Span | Servei | Inici (ms) | Durada (ms) | Atributs rellevants / observacions |
|---|---|---|---|---|
kong.proxy POST /api/v1/comandes |
Kong | 0 | 482 | http.status_code=201, km0.request_id=c1d2... |
⎿ POST /comandes |
comandes-2 | 3 | 476 | km0.comanda_id=P-2026-000125, enduser.pseudo=u:9f1c3a |
⎿ saga.reservar_estoc |
comandes-2 | 6 | 118 | span manual de saga_comanda.py |
⎿ km0.Inventari/ReservarEstoc (client) |
comandes-2 | 7 | 116 | rpc.grpc.status_code=0, deadline 300 ms |
⎿ km0.Inventari/ReservarEstoc (servidor) |
inventari-1 | 9 | 111 | net.peer.name=comandes-2 (mTLS, SPIFFE) |
⎿ SELECT km0_inventari |
inventari-1 | 11 | 4 | db.statement=SELECT ... FROM estoc WHERE producte=$1 FOR UPDATE |
⎿ UPDATE km0_inventari |
inventari-1 | 16 | 98 | db.statement=UPDATE estoc SET ...: espera de bloqueig |
⎿ COMMIT |
inventari-1 | 115 | 5 | |
⎿ saga.cobrar |
comandes-2 | 126 | 331 | |
⎿ km0.Pagaments/Cobrar (client → servidor) |
comandes-2 → pagaments | 127 | 329 | |
⎿ POST passarella.example/charges |
pagaments | 131 | 318 | http.status_code=200: la passarel·la externa domina la latència |
⎿ INSERT km0_comandes (Cassandra) |
comandes-2 | 459 | 9 | db.system=cassandra, consistència LOCAL_QUORUM |
⎿ publicar comandes.esdeveniments |
comandes-2 | 469 | 8 | messaging.destination=comandes.esdeveniments |
⎿ consumir comandes.esdeveniments |
analitica | 612 | 21 | fill de "publicar": comença 140 ms després que l'Anna ja tingui el seu 201 |
El que la traça explica i les mètriques no podien:
- Dels 482 ms, 318 són la passarel·la de pagaments externa. Res a
km0/no els pot accelerar; sí que es pot paral·lelitzar la reserva i el cobrament, o acceptar la comanda "pendent de confirmar" (07-04). - L'
UPDATEd'inventariva trigar 98 ms en una operació que normalment en triga 2: hi ha contenció de bloqueig sobrevi-crianca(la Setmana de la Verema, moltes comandes del mateix producte). La mètrica p99 deReservarEstocla mostrava pujant; la traça diu en quina sentència. - El span d'
analiticapertany a la mateixa traça encara que passi després de la resposta: és l'efecte de propagartraceparenta les capçaleres Kafka. - Cada span porta el
trace_id, i per tant un clic porta a les línies de log dels cinc serveis amb aquest identificador, en ordre.
- Unir els tres senyals
L'observabilitat no són tres eines, sinó un recorregut: la mètrica diu que alguna cosa va malament, la traça diu on, el log diu per què. Grafana permet que aquest recorregut siguin tres clics:
- Exemplars: Prometheus pot desar, al costat de cada bucket de l'histograma, el
trace_idd'una observació recent. Aprometheus_clientes passaexemplar={"trace_id": ...}aobserve(); al gràfic del p99, cada punt té un rombe que obre aquesta traça. És l'enllaç mètrica → traça. - Enllaç log ↔ traça: el datasource de Loki a Grafana es configura amb una derived field que reconeix
trace_idal JSON i mostra un botó "veure traça a Tempo"; i Tempo, a la inversa, amb "logs d'aquesta traça" que llança{servei=~".+"} | json | trace_id="..."a Loki. - Traça → mètrica: Tempo pot generar mètriques RED a partir dels spans (span metrics), útil per a serveis sense instrumentar amb Prometheus.
# km0/observabilitat/grafana/provisioning/datasources/observabilitat.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: "servei"}]
tracesToMetrics:
datasourceUid: prometheusUn exemple d'exemplar a metriques.py (07-01), ara que existeix el context de traça:
# km0/serveis/comu/metriques.py (afegit)
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
# dins de mesurar(), al finally:
PETICIO_DURADA.labels(...).observe(durada, exemplar=_exemplar_actual())Amb això, el dissabte a les 10:13 el recorregut d'en Jordi és: alerta ComandesBurnRateRapid → panell amb la taxa d'error → rombe d'exemplar sobre el pic → traça d'una comanda fallida amb el span ReservarEstoc en vermell i rpc.grpc.status_code=UNAVAILABLE → "logs d'aquesta traça" → {"esdeveniment": "circuit_obert", "dependencia": "inventari"}... que és on entra la resta del mòdul.
Errors Comuns i Consells
- Logs de text lliure "perquè es llegeixen millor". Es llegeixen millor un a un; no es consulten. Nom d'esdeveniment estable + camps; per llegir en local,
structlogté un renderitzador de consola amb colors. - Promocionar
trace_idocomanda_ida etiqueta de Loki. Cada valor crea un stream i Loki es degrada igual que Prometheus amb cardinalitat alta. Etiquetes:servei,instancia,nivell,entorn. La resta, al JSON. - Registrar capçaleres o cossos complets en
debugi deixar-ho activat. És la fuita de tokens i dades personals més habitual. Un processor de censura com a xarxa, idebugapagat per defecte. - Perdre el context als salts asíncrons. Un
ThreadPoolExecutor, unasyncio.create_tasko un missatge Kafka sensepropagate.injecttrenquen la traça. Davant d'una traça que "s'acaba de sobte", busca el salt sense propagació. - Mostreig només en capçalera. Descarta el 90 % dels errors. Tail sampling al Collector per quedar-se amb errors i traces lentes.
- Traces sense atributs de negoci. Un span
POST /comandessensekm0.comanda_idobliga a buscar per temps. Afegeix l'identificador de domini com a atribut (a traces i logs no hi ha problema de cardinalitat). - Confondre logs operatius amb auditoria. Retenció, immutabilitat i accés són diferents (06-05). Ni l'auditoria va a Loki amb 14 dies, ni els logs de depuració van al bucket amb object lock.
- Enviar traces directament al backend des de cada servei. Acobla tots els serveis al backend i fa impossible el tail sampling. Un Collector al mig.
Exercicis
Exercici 1. En Marc reporta a suport que la seva comanda de la Setmana del Formatge Artesà "ha donat error" i aporta l'X-Request-Id que la web li ha mostrat: 7a1b.... La Marta obre Grafana. (a) Escriu la consulta LogQL que reuneix tot el que va passar amb aquesta petició a tots els serveis. (b) La consulta retorna línies de Kong i comandes amb trace_id 9c4d..., però cap d'inventari; en canvi, hi ha línies d'inventari amb un trace_id diferent al mateix segon. Què s'ha trencat i on ho buscaries? (c) La Marta troba a la línia de comandes el camp "usuari": "u:2b77e1". Com confirma que és en Marc sense que ningú més no ho pugui fer pel seu compte, i per què el log no conté el seu email?
Exercici 2. L'equip decideix que repartiment publiqui una línia de log per cada posició rebuda (140 repartidors, una cada 5 s) amb el nivell info, i que debug inclogui el cos del missatge. (a) Calcula quantes línies al dia genera només repartiment i estima el volum si cada línia ocupa 350 bytes. (b) Proposa una alternativa que conservi la capacitat d'investigar "per què el panell mostrava furgoneta-3-017 a Girona quan era a Lleida" sense aquest volum. (c) Quin problema de privadesa afegeix el debug amb el cos, i amb quin mecanisme d'aquesta lliçó el mitigaries com a últim recurs?
Exercici 3. La traça de la secció 8 mostra 318 ms a la passarel·la de pagaments i 98 ms d'espera de bloqueig a inventari. Un company proposa "posar TraceIdRatioBased(1.0) a tots els serveis per no perdre cap traça com aquesta". (a) Calcula quants spans per segon rebria el Collector amb 100 peticions/s a comandes i 12 spans per traça, més el trànsit de cataleg (400 peticions/s, 4 spans). (b) Explica per què la configuració d'otel-collector.yaml de la secció 7 ja hauria conservat aquesta traça encara que el ratio fos 0,1, i què hauria passat si el ratio 0,1 s'hagués configurat sense ParentBased. (c) Un span de la compensació de la saga, executat dues hores després que la comanda fallés, apareix com a traça nova i no dins de la de la comanda. Què hi falta i on?
Solucions
Exercici 1.
(a) {servei=~".+"} | json | request_id="7a1b...", amb el rang de temps acotat al moment de la comanda (Loki escaneja el contingut, així que el rang importa). (b) Kong i comandes comparteixen trace_id, per tant la propagació HTTP funciona; inventari té un altre trace_id al mateix instant, per tant inventari està creant traces noves en lloc de continuar la de comandes: el traceparent no arriba a les metadades gRPC o no es llegeix. Sospitosos per ordre: GrpcInstrumentorClient no instrumentat a comandes (el client no injecta), un interceptor propi a inventari que reconstrueix les metadades i descarta traceparent (revisar AuthInterceptor/ServeiInterceptor de 06-04), o un stub gRPC creat abans de cridar configurar_traces. Es comprova amb un log.debug de les metadades rebudes a inventari (censurant authorization) o mirant a Tempo si el span servidor d'inventari és arrel. (c) La pseudonimització de 06-02 és una funció amb clau (HMAC) sobre l'identificador de l'usuari: la Marta demana al servei d'identitats, amb el seu rol d'operador i el seu tiquet de suport (queda auditat, 06-05), que resolgui u:2b77e1; o calcula el pseudònim del sub d'en Marc i compara. El log no conté l'email perquè és una dada personal que viatja a Loki, a MinIO i a la pantalla de qualsevol amb accés a Grafana: es registra el pseudònim, que no identifica sense la clau, i el comanda_id, que és el que suport necessita.
Exercici 2.
(a) 140 × (86 400 / 5) = 2 419 200 línies al dia; a 350 bytes, uns 847 MB diaris sense comprimir només de posicions, uns 25 GB al mes. (b) Registrar en info només allò anòmal (posicions rebutjades per validació, salts impossibles entre dues posicions, missatges fora d'ordre) i un resum per repartidor i minut (posicions_rebudes, ultima_ts); per reconstruir la ruta de furgoneta-3-017, les dades ja són a repartiment.posicions amb retenció d'hores i al llac HDFS (04-02): investigar consultant aquestes dades, no el log. A més, una traça mostrejada per repartidor cada N missatges dona el recorregut intern sense registrar-los tots. (c) El cos conté la posició exacta d'un empleat amb el seu identificador, una dada personal (06-05, exercici 3); si algú activa debug en producció, es copia a Loki sense les garanties d'accés de repartiment.posicions. Com a últim recurs, el processor censurar amb posicio/lat/lon a CLAUS_PROHIBIDES (i que el cos es registri com a camps, mai com a cadena opaca, perquè el censor pugui actuar); però la política correcta és no registrar el cos.
Exercici 3.
(a) comandes: 100 × 12 = 1 200 spans/s; cataleg: 400 × 4 = 1 600 spans/s; total uns 2 800 spans/s, 242 milions al dia, que el Collector ha de rebre i retenir 10 s cadascun per decidir (uns 28 000 spans en memòria de manera permanent, factible), i dels quals Tempo escriuria el que les polítiques deixin passar. Amb TraceIdRatioBased(1.0) el cost és a la xarxa i al Collector, no a l'emmagatzematge, sempre que hi hagi tail sampling. (b) La traça dura 482 ms i la política lentes conserva tot el que superi 500 ms... no la conservaria per durada; però tampoc no importa: si va tenir un span amb error, la política errors la desa; si no, cau en el 10 % probabilístic. Amb el ratio 0,1 en capçalera, la decisió la va prendre Kong (o comandes) en crear la traça, i ParentBased fa que tots els serveis la respectin: o es mostreja sencera o cap no la genera. Sense ParentBased, cada servei tiraria el seu propi dau: comandes la conservaria i inventari no, i la traça quedaria amb forats (el span ReservarEstoc servidor absent, precisament el dels 98 ms). (c) Falta desar el traceparent de la comanda a la fila de la taula sagas quan es crea la saga, i en executar la compensació reconstruir el context amb propagate.extract sobre aquest valor i crear el span de compensació amb context=ctx (o com a link a la traça original si es prefereix una traça nova enllaçada). És a saga_comanda.py, al punt on el relay o el planificador reprenen una saga pendent.
Conclusió
Amb aquesta lliçó l'observabilitat de Quilòmetre Zero té els seus tres senyals. Els logs deixen de ser text en fitxers de contenidors efímers per ser esdeveniments JSON amb camps comuns (ts en UTC, nivell, servei, instancia, trace_id, request_id, usuari pseudonimitzat, esdeveniment), sense secrets ni dades personals, recollits per Promtail a cada node i emmagatzemats a Loki amb etiquetes de baixa cardinalitat, consultables amb LogQL per l'X-Request-Id que Kong assigna o pel trace_id que OpenTelemetry injecta. Les traces donen el recorregut: spans imbricats amb traceparent W3C propagat per HTTP i gRPC sense codi, i a mà per les capçaleres de Kafka i per la taula sagas, enviats a un Collector que mostreja en cua per conservar errors i lentitud, i emmagatzemats a Tempo. I els tres s'uneixen a Grafana amb exemplars i enllaços log↔traça, de manera que una alerta es converteix en tres clics en una traça amb el span culpable i les línies de log que l'expliquen.
Ja sabem que falla, on i per què. El recorregut d'en Jordi acabava en un log de comandes que deia circuit_obert cap a inventari, i a Tempo un span ReservarEstoc en vermell amb UNAVAILABLE: inventari-1 no respon. La pregunta següent no és d'observabilitat, sinó de supervivència: com detecta la plataforma que un node ha caigut i no simplement va lent? Qui decideix que inv-vlc passi a ser primari, i com s'evita que inv-bcn torni creient que encara ho és? I si el que s'ha perdut no és un procés sinó les dades de km0_inventari de les darreres dues hores? La lliçó següent tracta la gestió de fallades i la recuperació: detecció, failover amb tancat, còpies de seguretat i PITR, i què fer en els minuts en què tot això està passant.
Curs d'Arquitectures Distribuïdes
Mòdul 1: Introducció als Sistemes Distribuïts
- Conceptes Bàsics de Sistemes Distribuïts
- Models de Sistemes Distribuïts
- Avantatges i Desafiaments dels Sistemes Distribuïts
- Les Fal·làcies de la Computació Distribuïda
- Temps, Rellotges i Ordenació d'Esdeveniments
- Del Monòlit a la Plataforma Distribuïda: el Cas Quilòmetre Zero
Mòdul 2: Comunicació en Sistemes Distribuïts
- Protocols de Comunicació
- RPC i RMI
- gRPC i Serialització de Dades
- Missatgeria i Cues de Missatges
- Patrons de Comunicació Asíncrona
Mòdul 3: Consistència i Replicació
- Models de Consistència
- El Teorema CAP i PACELC
- Algorismes de Consens
- Replicació de Dades
- Transaccions Distribuïdes i Sagues
Mòdul 4: Emmagatzematge Distribuït
- Particionament de Dades i Hashing Consistent
- Sistemes de Fitxers Distribuïts
- Emmagatzematge d'Objectes
- Bases de Dades Distribuïdes
- Memòries Cau Distribuïdes
Mòdul 5: Computació Distribuïda
- Models de Computació Distribuïda
- MapReduce i Hadoop
- Spark i Computació en Memòria
- Processament de Fluxos de Dades
- Planificació de Treballs i Pipelines de Dades
Mòdul 6: Seguretat en Sistemes Distribuïts
- Autenticació i Autorització
- Xifratge i Protecció de Dades
- Gestió d'Identitats
- Seguretat entre Serveis: mTLS i Gestió de Secrets
- Passarel·les d'API, Limitació de Taxa i Auditoria
Mòdul 7: Monitoratge i Manteniment
- Monitoratge de Sistemes Distribuïts
- Logs Centralitzats i Traçabilitat Distribuïda
- Gestió de Fallades i Recuperació
- Patrons de Resiliència: Timeouts, Reintents i Circuit Breaker
- Automatització i Orquestració
- Proves en Sistemes Distribuïts i Enginyeria del Caos
