сохранять непрерывность распределённых трейсов на каждой границе сервисов
Передача контекста между микросервисами: практические паттерны 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.
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)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 активного спана, чтобы с дашборда переходить от строки лога к трейсу.
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), чтобы бэкенды корректно группировали асинхронные прыжки.
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 с исходными контекстами — бэкенд покажет причинную связь без единого родителя на все элементы.
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)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.
