сохранять непрерывность распределённых трейсов на каждой границе сервисов

Передача контекста между микросервисами: практические паттерны OpenTelemetry

14 минут

Оборванные трейсы чаще всего не из‑за отсутствия инструментирования, а из‑за потерянного traceparent на HTTP-вызове, в заголовке Kafka, в метаданных gRPC или в фоновом потоке. Освойте инъекцию, транспорт и извлечение — и весь путь запроса останется в одном трейсе.

Почему трейсы рвутся, даже когда сервисы инструментированы

Пользовательский запрос через HTTP API, gRPC и очереди даёт полезный водопад только если один общий идентификатор трейса переживает каждый прыжок. Команды ставят SDK OpenTelemetry и всё равно видят сиротские корни: сервис B начинает новый трейс, потому что A не отправил заголовок W3C traceparent; потребитель Kafka не видит заголовков продюсера; cron-воркер собирает контекст с нуля; метаданные gRPC не подключены; батч из ста сообщений схлопывается в один спан без родителя. Каждый разрыв поднимает MTTR — дежурный сверяет метки времени вместо одного непрерывного трейса. Инструментирование без передачи контекста — шум. Здесь — паттерны на границах; топология коллекторов и автоинъекция в Kubernetes — в соседних гайдах по трассированию.

Модель: инъекция, транспорт, извлечение

Передача контекста в OpenTelemetry всегда из трёх шагов. Инъекция записывает активный контекст спана в носитель — HTTP-заголовки, метаданные gRPC или заголовки сообщений — обычно поля W3C Trace Context traceparent и tracestate, при необходимости плюс baggage. Транспорт — протокол, который должен донести эти поля без потерь. Извлечение восстанавливает Context OpenTelemetry у получателя, чтобы следующий спан стал дочерним к вышестоящему. Предпочитайте автоинструментирование фреймворков; API propagate оставляйте для самописных клиентов. Настройте пропагаторы один раз при старте через OTEL_PROPAGATORS или set_global_textmap — по умолчанию уже идут tracecontext и baggage.

Python · инъекция и извлечение на HTTP
from opentelemetry import trace
from opentelemetry.propagate import extract, inject

tracer = trace.get_tracer("order-service")

def call_downstream(client):
    with tracer.start_as_current_span("call-inventory") as span:
        span.set_attribute("peer.service", "inventory")
        headers = {}
        inject(headers)  # writes traceparent, tracestate, baggage
        return client.get("http://inventory/api/stock", headers=headers)

def handle_incoming(request):
    ctx = extract(dict(request.headers))
    with tracer.start_as_current_span("handle-request", context=ctx) as span:
        span.set_attribute("http.route", "/api/orders")
        return call_downstream(http_client)
Python · совместная работа W3C и B3 при миграции
from opentelemetry.propagate import set_global_textmap
from opentelemetry.propagators.b3 import B3MultiFormat
from opentelemetry.propagators.composite import CompositePropagator
from opentelemetry.baggage.propagation import W3CBaggagePropagator
from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator

# Prefer OTEL_PROPAGATORS=tracecontext,baggage,b3multi in production.
set_global_textmap(
    CompositePropagator(
        [
            TraceContextTextMapPropagator(),
            W3CBaggagePropagator(),
            B3MultiFormat(),
        ]
    )
)

HTTP и gRPC: сначала перехватчики фреймворка

Для HTTP используйте автоинструментирование языка — FastAPI, Flask, requests, httpx, обёртки Go net/http — чтобы каждый исходящий клиент и входящий сервер наследовали контекст. Ручная инъекция — запасной путь для самописных клиентов. Для gRPC подключите инструментирование клиента и сервера, чтобы метаданные несли те же поля W3C; не изобретайте параллельную схему заголовков. Идентификаторы тенанта или региона кладите в baggage или прикладные метаданные осознанно, а не вместо traceparent. В структурированные логи пишите trace_id активного спана, чтобы с дашборда переходить от строки лога к трейсу.

Python · инструментирование клиента и сервера gRPC
from concurrent import futures

import grpc
from opentelemetry.instrumentation.grpc import (
    GrpcInstrumentorClient,
    GrpcInstrumentorServer,
)

GrpcInstrumentorClient().instrument()
GrpcInstrumentorServer().instrument()

channel = grpc.insecure_channel("inventory:50051")
stub = InventoryStub(channel)

server = grpc.server(futures.ThreadPoolExecutor(max_workers=10))
# Register servicers, then start — interceptors are already active.

Kafka и другие асинхронные границы

Очереди — самая частая тихая поломка. Продюсер обязан положить контекст в заголовки сообщения; потребитель — извлечь его до бизнес-спанов. KafkaInstrumentor оборачивает продюсеры и консюмеры kafka-python и создаёт спаны PRODUCER и CONSUMER. Итератор консюмера всё равно может потерять контекст в пользовательском коде — извлекайте заголовки сами и передавайте context= в start_as_current_span для бизнес-спана. Тот же приём для свойств RabbitMQ, атрибутов SQS и заголовков NATS: один словарь-носитель, inject при публикации, extract при чтении. Используйте семантические соглашения messaging (messaging.system, messaging.destination.name), чтобы бэкенды корректно группировали асинхронные прыжки.

Python · публикация и потребление Kafka с контекстом
import json

from kafka import KafkaConsumer, KafkaProducer
from opentelemetry import trace
from opentelemetry.instrumentation.kafka import KafkaInstrumentor
from opentelemetry.propagate import extract, inject

KafkaInstrumentor().instrument()
tracer = trace.get_tracer("order-service")
producer = KafkaProducer(bootstrap_servers="kafka:9092")

def publish_order(order):
    with tracer.start_as_current_span("publish-order") as span:
        span.set_attribute("order.id", order["id"])
        span.set_attribute("messaging.system", "kafka")
        span.set_attribute("messaging.destination.name", "orders")
        headers = {}
        inject(headers)
        producer.send(
            "orders",
            key=order["id"].encode(),
            value=json.dumps(order).encode(),
            headers=[(k, v.encode()) for k, v in headers.items()],
        )

def consume_orders(consumer: KafkaConsumer):
    for msg in consumer:
        carrier = {
            key: (value.decode() if isinstance(value, (bytes, bytearray)) else value)
            for key, value in (msg.headers or [])
        }
        ctx = extract(carrier)
        with tracer.start_as_current_span("consume-order", context=ctx) as span:
            order = json.loads(msg.value)
            span.set_attribute("order.id", order["id"])
            reserve_inventory(order)

Baggage, потоки и span links для батчей

Baggage переносит прикладные ключи вроде tenant.id или feature.flag вместе с трейсом. Используйте умеренно — каждый прыжок платит размером заголовков, а секреты в baggage класть нельзя. Перед запуском потока или задачи в executor скопируйте контекст и восстановите его внутри воркера, иначе дочерний спан станет новым корнем. При веерной обработке батча иерархия parent-child часто неверна: создавайте независимые спаны обработки и связывайте их через span links с исходными контекстами — бэкенд покажет причинную связь без единого родителя на все элементы.

Python · baggage и копирование контекста в поток
import threading

from opentelemetry import baggage, context, trace

tracer = trace.get_tracer("checkout")

def enqueue_with_tenant(tenant_id: str):
    ctx = baggage.set_baggage("tenant.id", tenant_id)
    token = context.attach(ctx)
    try:
        with tracer.start_as_current_span("enqueue"):
            captured = context.get_current()

            def worker():
                restore = context.attach(captured)
                try:
                    with tracer.start_as_current_span("background-work"):
                        assert baggage.get_baggage("tenant.id") == tenant_id
                        do_work()
                finally:
                    context.detach(restore)

            threading.Thread(target=worker).start()
    finally:
        context.detach(token)
Python · span links для веерной обработки батча
from opentelemetry import trace
from opentelemetry.propagate import extract

tracer = trace.get_tracer("batch-worker")

def process_batch(messages):
    with tracer.start_as_current_span("batch-process") as batch_span:
        batch_span.set_attribute("messaging.batch.message_count", len(messages))
        for message in messages:
            carrier = dict(message.headers)
            extracted = extract(carrier)
            linked = trace.get_current_span(extracted).get_span_context()
            with tracer.start_as_current_span(
                "process-item",
                links=[trace.Link(linked)],
            ) as item_span:
                item_span.set_attribute("messaging.message.id", message.id)
                handle(message)

Проверяйте передачу тестами и сигналами покрытия

Добавьте интеграционную проверку: клиент делает inject, нижестоящий сервис продолжает тот же trace_id — сравнивайте отформатированные hex-строки, а не сырое целое с самодельным заголовком ответа, если вы сами не эхоите format_trace_id. В бэкенде алертируйте на рост доли корневых спанов у сервисов, которые почти всегда должны иметь родителя (шлюзы — исключение). Метрики пайплайна вроде ошибок экспорта тоже важны, но доля сиротских корней — самый прямой сигнал здоровья передачи. Держите head sampling parent-based, чтобы выбранный корень тянул дочерние спаны; иначе обрыв посередине выглядит как баг пропагации.

Операционный чеклист непрерывных трейсов

Стандартизируйте W3C Trace Context везде; B3 оставляйте только на время миграции. Инструментируйте библиотеки на краю фреймворка до ручных вызовов inject. Проверьте каждый асинхронный продюсер и потребитель на заголовки. Запретите PII в baggage и атрибутах спанов. Связывайте логи с trace_id и span_id. После появления новой очереди или gRPC-клиента пересматривайте сиротские корни. Передача контекста — невидимый контракт распределённого трассирования: чините inject и extract на каждой границе — и следующий кросс-сервисный инцидент станет одним водопадом, а не археологией.

Коллекторы, сэмплирование и автоинструментирование в кластере разобраны в гайде по распределённому трассированию OpenTelemetry в Kubernetes.

Непрерывные трейсы дальше проходят через гайд по production-pipeline OpenTelemetry Collector.