Допустим, поток должен передать данные в asyncio.Queue, обслуживаемую другим потоком с event loop. Можно ли безопасно вызвать queue.put из исходного потока?
Нет. Методы asyncio.Queue не предназначены для прямого вызова из другого потока: очередь не является потокобезопасной. Поток должен передать операцию в поток event loop через loop.call_soon_threadsafe или asyncio.run_coroutine_threadsafe.
asyncio проектировалась вокруг однопоточного event loop с кооперативным выполнением корутин. Такой подход уменьшает необходимость в блокировках и делает взаимодействие задач дешевле, но предполагает, что asyncio-примитивы обслуживаются своим loop-потоком.
Для связи с внешними потоками предусмотрены специальные механизмы потокобезопасной постановки callback или корутины в event loop. Они отделяют синхронизацию между потоками от обычных операций очереди.
Прямой вызов put из другого потока некорректен по двум причинам. Во-первых, put является корутиной, поэтому без правильной передачи в event loop она не будет выполнена. Во-вторых, даже put_nowait нельзя считать универсально безопасным для конкурентного вызова из другого потока.
Ошибочное решение может привести к потерянным уведомлениям, некорректной работе ожиданий или трудно воспроизводимым зависаниям. GIL не превращает asyncio.Queue в потокобезопасную структуру: потокобезопасность требует согласованного взаимодействия с event loop, а не только защиты отдельных операций интерпретатора.
Если операция не должна ждать свободного места в очереди, поток может безопасно запланировать callback в loop через call_soon_threadsafe. Сам queue.put_nowait тогда выполнится уже в правильном потоке.
Если очередь ограничена и операция должна дождаться свободного места, используют run_coroutine_threadsafe. Он принимает корутину, планирует её в указанном loop и возвращает объект concurrent.futures.Future; поток при необходимости может дождаться результата через result().
call_soon_threadsafe подходит для короткой неблокирующей постановки callback. run_coroutine_threadsafe лучше отражает асинхронную семантику await queue.put(...), но ожидание его результата из потока может заблокировать этот поток и создать взаимное ожидание, если loop зависит от него.
Важно передавать именно тот event loop, который реально запущен и владеет очередью. После остановки loop постановка новой операции завершится ошибкой или не даст ожидаемого результата. Для сложных двусторонних мостов между потоками и asyncio можно использовать специализированную потокобезопасную очередь, но обычная asyncio.Queue сама по себе таким мостом не является.
Синхронный драйвер получает события в отдельном потоке, а веб-сервис обрабатывает их в asyncio. Разработчик сначала вызывает queue.put_nowait прямо из callback драйвера. На тестовой нагрузке это работает, но при заполненной очереди и одновременной остановке сервиса появляются потерянные события.
Вариант с общей блокировкой вокруг очереди не решает проблему: блокировка не делает внутренние ожидания asyncio корректными и может задержать event loop. Вариант с постоянным опросом очереди из asyncio проще, но добавляет задержку и расходует ресурсы.
Выбранное решение — передавать события через call_soon_threadsafe, а для ограниченной очереди — через run_coroutine_threadsafe. При завершении сервиса сначала прекращают приём новых событий, затем корректно завершают loop и проверяют результаты отправленных операций. Это сохраняет владение asyncio-очередью за одним loop-потоком и позволяет явно контролировать переполнение.
1. Можно ли вызвать queue.put_nowait из другого потока, если операция кажется атомарной?
Нет, гарантии потокобезопасности у asyncio.Queue нет. Даже отсутствие await внутри конкретного вызова не означает, что его можно безопасно выполнять вне loop-потока. Операцию следует запланировать через call_soon_threadsafe.
2. Чем отличается run_coroutine_threadsafe от asyncio.create_task в этой ситуации?
create_task предназначен для вызова из потока, в котором работает соответствующий event loop, и создаёт задачу в текущем loop-контексте. run_coroutine_threadsafe специально принимает loop и позволяет другому потоку передать ему корутину, возвращая синхронный concurrent.futures.Future для наблюдения за результатом.
3. Что произойдёт, если поток вызовет future.result(), а корутина в loop ждёт этот же поток?
Возникает взаимная блокировка. Поток ждёт завершения корутины, а корутина не может продолжить работу, пока поток не освободится. Поэтому result() из внешнего потока следует вызывать с осторожностью, часто с тайм-аутом, либо строить протокол обмена без синхронного ожидания.