ГлавнаяБлогCorrelation ID в asyncio: как связать логи трёх сервисов
Python

Correlation ID в asyncio: как связать логи трёх сервисов

Узнайте, как с помощью contextvars и asyncio.to_thread добавить сквозной correlation ID в логи асинхронных Python-сервисов. Практические примеры и тесты.

Al
Редакция Algolitalgolit.ru
8 мин чтения4 августа 2026 г.

Почему логи трёх сервисов не сходятся

Заказ падает в 14:02. В тикете поддержки цитата клиента: сначала ошибка оплаты, потом повторная попытка, снова ошибка. Три лог-строки с одним временным окном — и ни одна не объясняет, что произошло. Платёжный сервис записал payment.authorized в 14:02:11, сервис заказов — order.failed в 14:02:11, а воркер, который должен был их согласовать, — worker.idle в 14:02:12. Три сервиса, три потока логов, три разные истории. Между этими строками потеряно изменение состояния, и никто не может сказать, где именно, потому что ничто в логах их не связывает.

Это классический симптом системы, где каждая лог-строка локально правдива, но глобально бесполезна. Решение — correlation ID: одна непрозрачная строка, которая путешествует с каждым запросом, пересекает границы сервисов и добавляется в каждую лог-строку. В синхронном коде это решённая задача — заголовок, middleware, готово. В асинхронном Python интересно то, что механизм, к которому вы потянетесь первым, тихо неверен, а рабочий спрятан в углу стандартной библиотеки, который большинство людей никогда не импортировали.

Почему thread-local состояние ломается в асинхронном коде

Инстинктивное решение — threading.local(). Один поток, один запрос, один контекст — эта ментальная модель сделала thread-local storage популярным, и именно её нарушает асинхронный код.

Цикл событий asyncio выполняет тысячи задач в одном потоке. Поток не меняется, когда приложение переключается с обработки запроса A на запрос B; меняется задача. threading.local() хранит данные по идентификатору потока, поэтому все задачи в одном потоке читают и пишут в один и тот же слот. Запрос A устанавливает request_id = "a1", ожидает I/O, и пока он приостановлен, запрос B устанавливает request_id = "b2". Когда A возобновляется и пишет следующую строку лога, фильтр читает request_id и помечает строку как "b2".

import threading
local = threading.local()

async def handle(request_id: str) -> None:
    local.request_id = request_id
    await some_io()  # цикл переключается на другую задачу здесь
    print(local.request_id)  # может быть уже ID другого запроса

Результат — не отсутствие ID, а хуже: неправильный ID. Лог-строки помечаются ID соседнего запроса, и восстанавливаемый таймлайн становится активно вводящим в заблуждение. Это чистый пример провала «локально правда, глобально бесполезно». Поэтому стандартный совет для асинхронного кода — не хранить контекст в местах, которые рантайм может подменить под вами.

Contextvars: состояние, которое путешествует с задачей

contextvars существует именно для этого. Контекстная переменная хранит одно значение на контекст выполнения, и каждая задача asyncio несёт свой собственный контекст. Когда цикл приостанавливает одну задачу и возобновляет другую, каждая задача видит свою копию каждой контекстной переменной — без общего слота и без перекрёстных помех.

import contextvars

request_id: contextvars.ContextVar[str] = contextvars.ContextVar("request_id", default="")

async def handle(request_id_value: str) -> None:
    request_id.set(request_id_value)
    await some_io()  # цикл может свободно переключать задачи
    print(request_id.get())  # всё ещё "a1" для этой задачи

Два свойства делают это безопасным. Первое: присваивание внутри задачи видно только этой задаче и тому, что она ожидает (await); соседние задачи его не наблюдают. Второе: когда задача завершается, её контекст отбрасывается, поэтому нет утечки в следующий запрос. Значение следует за логической единицей работы, а не за потоком — именно такая семантика нужна correlation ID.

Стоит назвать цену: contextvars не бесплатен. Установка и чтение контекстной переменной включают поиск по словарю в текущем контексте, и каждый await может вызывать служебные операции с контекстом. В горячем пути логирования стоимость ничтожна по сравнению с уже выполняемым I/O, но если вы планируете использовать контекстную переменную внутри тесного численного цикла — сначала измерьте.

Middleware для correlation ID в тридцать строк

У middleware три задачи: принять ID от вышестоящего сервиса, если он есть; сгенерировать новый, если его нет; и вернуть ID в ответе, чтобы вызывающий мог прикрепить его к своим логам. С contextvars реализация короткая — можно прочитать за один проход.

import contextvars
import uuid
from starlette.middleware.base import BaseHTTPMiddleware

request_id_var: contextvars.ContextVar[str] = contextvars.ContextVar("request_id", default="")

class CorrelationMiddleware(BaseHTTPMiddleware):
    async def dispatch(self, request, call_next):
        incoming = request.headers.get("X-Request-ID", "").strip()
        rid = incoming or uuid.uuid4().hex[:16]
        request_id_var.set(rid)
        response = await call_next(request)
        response.headers["X-Request-ID"] = rid
        return response

Та же форма переносится на middleware aiohttp, зависимости FastAPI и обработчики сырого asyncio. Важна не привязка к фреймворку, а то, что ID живёт в request_id_var, поэтому любая функция, которая вызывается внутри этого запроса и делает await, может прочитать его без передачи параметра через десятки мест вызова.

Это важнее, чем кажется. Альтернатива — передавать ID как явный аргумент в каждую функцию, которая логирует — превращает каждую сигнатуру в переносчик служебных данных. Добавление поля в датакласс, параметра в хелпер или аргумента в сторонний колбэк становится ломающим изменением. Контекстная переменная сохраняет поток данных неявным, а места вызова остаются честными: функция логирует то, что ей нужно, потому что контекст уже там.

Ловушка копирования контекста: потоки и экзекуторы

Вот угол, о который спотыкаются большинство реализаций: asyncio.create_task() копирует текущий контекст, но работа, отправленная в пул потоков через экзекутор, его не наследует.

Задача, вызывающая loop.run_in_executor(None, blocking_call), передаёт blocking_call рабочему потоку. Этот поток работает с контекстом по умолчанию, в котором request_id всё ещё пустой дефолт. Correlation ID молча исчезает именно на той границе, где он нужнее всего — на медленном блокирующем вызове, который занимает реальные секунды и пишет реальные ошибки.

import asyncio
import contextvars

request_id: contextvars.ContextVar[str] = contextvars.ContextVar("request_id", default="")

def blocking_call() -> None:
    print(request_id.get())  # "" — поток экзекутора имеет контекст по умолчанию

async def handle(rid: str) -> None:
    request_id.set(rid)
    await asyncio.get_running_loop().run_in_executor(None, blocking_call)

asyncio.run(handle("a1"))

Исправление — явно распространить контекст. Python 3.11 добавил asyncio.to_thread(), который копирует контекст вызывающего перед отправкой — замена в одну строку, которая позволяет ID пережить границу:

async def handle(rid: str) -> None:
    request_id.set(rid)
    await asyncio.to_thread(blocking_call)  # контекст копируется автоматически

Для кода, который должен остаться на run_in_executor, скопируйте контекст вручную и выполните вызываемый объект внутри него:

ctx = contextvars.copy_context()
await loop.run_in_executor(None, lambda: ctx.run(blocking_call))

Выберите одно правило и применяйте его везде: предпочитайте asyncio.to_thread и проверяйте каждый оставшийся run_in_executor на явный ctx.run. Correlation ID, который умирает в первом пуле потоков, хуже, чем отсутствие ID, потому что логи теперь выглядят полными, но молча теряют часть таймлайна, которая действительно важна.

Фильтр логирования: чтобы каждая строка несла ID

Middleware устанавливает переменную, но логи всё ещё должны её читать. Фильтр logging.Filter в Python выполняется для каждой записи перед её отправкой, что делает его естественным местом для добавления ID:

import logging

class RequestIDFilter(logging.Filter):
    def filter(self, record: logging.LogRecord) -> bool:
        record.request_id = request_id_var.get() or "-"
        return True

handler = logging.StreamHandler()
handler.addFilter(RequestIDFilter())
logging.basicConfig(
    handlers=[handler],
    level=logging.INFO,
    format="%(asctime)s %(request_id)s %(name)s %(message)s"
)

Каждая строка — из кода приложения, из логгеров библиотек, из путей предупреждений — теперь несёт ID, потому что фильтру всё равно, кто создал запись. Один фильтр, прикреплённый к корневому обработчику, покрывает и ваш код, и чужой.

Это свойство снова связывает три лога. Строка платежа, строка заказа и строка воркера несут одно и то же значение X-Request-ID, и тикет поддержки превращается в поиск: отфильтруйте агрегированный поток по этой строке — и весь таймлайн упавшего заказа станет непрерывным. Correlation ID не делает логи правдивыми; он делает их сортируемыми по правде.

Проверяем, что таймлайн восстановим

Реализация не закончена, пока не протестирован регресс перекрёстных помех. Два теста покрывают важные режимы отказа: параллельные задачи не должны делить ID, а работа в экзекуторе должна их наследовать.

import asyncio
import contextvars

def test_tasks_do_not_share_ids():
    seen: set[str] = set()

    async def worker(value: str) -> str:
        request_id.set(value)
        await asyncio.sleep(0)
        return request_id.get()

    async def main() -> None:
        results = await asyncio.gather(worker("a1"), worker("b2"))
        seen.update(results)

    asyncio.run(main())
    assert seen == {"a1", "b2"}  # ни одна задача не видела ID соседа

def test_executor_inherits_context():
    captured: list[str] = []

    def blocking() -> None:
        captured.append(request_id.get())

    async def main() -> None:
        request_id.set("a1")
        await asyncio.to_thread(blocking)

    asyncio.run(main())
    assert captured == ["a1"]

Запустите эти тесты в CI-задаче, которая выполняет ваши асинхронные тесты. Если какой-то из них падает, исправление — один из двух паттернов выше, а не новый конфигурационный флаг и не глобальная переменная.

Когда contextvars недостаточно

Контекстная переменная переносит ID через один процесс. Как только запрос пересекает сетевую границу, ID должен ехать в протоколе: заголовок X-Request-ID в HTTP, заголовок traceparent для OpenTelemetry, поле message_id в очереди. Middleware уже распространяет заголовок вперёд; тот же заголовок следует читать от вышестоящего сервиса и прикреплять к исходящим вызовам, чтобы цепочка из трёх сервисов давала один искомый ID, а не три локально уникальных.

Для систем, где уже работает OpenTelemetry, прагматичный шаг — подкрепить correlation ID контекстом трассировки, а не изобретать свой: читайте traceparent, сохраняйте trace ID в контекстную переменную и пусть фильтр его добавляет. Паттерн не меняется; меняется источник строки.

И честное ограничение: contextvars не заставит логи из упавшего процесса появиться. Он закрывает разрыв между лог-строками, которые были записаны; он не может воскресить строку, которая не была записана, потому что процесс умер во время сброса. Для этого всё ещё нужна надёжная буферизованная доставка. Но для обычного сбоя — того, где каждый сервис что-то записал и никто не мог соединить точки — correlation ID — это разница между таймлайном и кучей временных меток.

Что делать прямо сейчас

  1. Возьмите middleware из статьи и подключите его к вашему FastAPI или aiohttp приложению.
  2. Добавьте фильтр логирования RequestIDFilter в корневой обработчик.
  3. Замените все run_in_executor на asyncio.to_thread в критических путях.
  4. Добавьте два теста из раздела «Проверяем, что таймлайн восстановим» в ваш CI.
  5. Настройте передачу заголовка X-Request-ID между сервисами.

После этого следующий инцидент с рассинхронизацией логов станет поиском по одному ID, а не гаданием по трём разрозненным историям.

#correlation id#asyncio#contextvars#логирование#middleware
Al
Редакция Algolit

Пишем про алгоритмы, подготовку к собеседованиям и карьеру в IT — так, чтобы было понятно и полезно.

Хочешь закрепить знания на практике?

Решай задачи на Algolit — интерактивная платформа для обучения

Начать бесплатно →
Correlation ID в asyncio: как связать логи трёх сервисов | Algolit