Узнайте, как с помощью contextvars и asyncio.to_thread добавить сквозной correlation ID в логи асинхронных Python-сервисов. Практические примеры и тесты.
Заказ падает в 14:02. В тикете поддержки цитата клиента: сначала ошибка оплаты, потом повторная попытка, снова ошибка. Три лог-строки с одним временным окном — и ни одна не объясняет, что произошло. Платёжный сервис записал payment.authorized в 14:02:11, сервис заказов — order.failed в 14:02:11, а воркер, который должен был их согласовать, — worker.idle в 14:02:12. Три сервиса, три потока логов, три разные истории. Между этими строками потеряно изменение состояния, и никто не может сказать, где именно, потому что ничто в логах их не связывает.
Это классический симптом системы, где каждая лог-строка локально правдива, но глобально бесполезна. Решение — correlation ID: одна непрозрачная строка, которая путешествует с каждым запросом, пересекает границы сервисов и добавляется в каждую лог-строку. В синхронном коде это решённая задача — заголовок, middleware, готово. В асинхронном Python интересно то, что механизм, к которому вы потянетесь первым, тихо неверен, а рабочий спрятан в углу стандартной библиотеки, который большинство людей никогда не импортировали.
Инстинктивное решение — 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 существует именно для этого. Контекстная переменная хранит одно значение на контекст выполнения, и каждая задача 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 три задачи: принять 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, потому что логи теперь выглядят полными, но молча теряют часть таймлайна, которая действительно важна.
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-задаче, которая выполняет ваши асинхронные тесты. Если какой-то из них падает, исправление — один из двух паттернов выше, а не новый конфигурационный флаг и не глобальная переменная.
Контекстная переменная переносит ID через один процесс. Как только запрос пересекает сетевую границу, ID должен ехать в протоколе: заголовок X-Request-ID в HTTP, заголовок traceparent для OpenTelemetry, поле message_id в очереди. Middleware уже распространяет заголовок вперёд; тот же заголовок следует читать от вышестоящего сервиса и прикреплять к исходящим вызовам, чтобы цепочка из трёх сервисов давала один искомый ID, а не три локально уникальных.
Для систем, где уже работает OpenTelemetry, прагматичный шаг — подкрепить correlation ID контекстом трассировки, а не изобретать свой: читайте traceparent, сохраняйте trace ID в контекстную переменную и пусть фильтр его добавляет. Паттерн не меняется; меняется источник строки.
И честное ограничение: contextvars не заставит логи из упавшего процесса появиться. Он закрывает разрыв между лог-строками, которые были записаны; он не может воскресить строку, которая не была записана, потому что процесс умер во время сброса. Для этого всё ещё нужна надёжная буферизованная доставка. Но для обычного сбоя — того, где каждый сервис что-то записал и никто не мог соединить точки — correlation ID — это разница между таймлайном и кучей временных меток.
RequestIDFilter в корневой обработчик.run_in_executor на asyncio.to_thread в критических путях.X-Request-ID между сервисами.После этого следующий инцидент с рассинхронизацией логов станет поиском по одному ID, а не гаданием по трём разрозненным историям.
Хочешь закрепить знания на практике?
Решай задачи на Algolit — интерактивная платформа для обучения
Начать бесплатно →