Skip to content

Instantly share code, notes, and snippets.

@artemonsh
Created August 26, 2026 09:31
Show Gist options
  • Select an option

  • Save artemonsh/dfb487bf0bba4865a9ea4c686badcf83 to your computer and use it in GitHub Desktop.

Select an option

Save artemonsh/dfb487bf0bba4865a9ea4c686badcf83 to your computer and use it in GitHub Desktop.
Пример запуска асинхронных функций в Celery через пул потоков
import asyncio
import threading
from collections.abc import Coroutine
from typing import Any
class CeleryEventLoop:
def __init__(self) -> None:
self._loop: asyncio.AbstractEventLoop | None = None
self._lock = threading.Lock()
def get_loop(self) -> asyncio.AbstractEventLoop:
if self._loop is not None and self._loop.is_running():
return self._loop
with self._lock:
if self._loop is not None and self._loop.is_running():
return self._loop
self._loop = asyncio.new_event_loop()
thread = threading.Thread(
target=self._run_loop,
args=(self._loop,),
name="celery-event-loop",
daemon=True,
)
thread.start()
return self._loop
def _run_loop(self, loop: asyncio.AbstractEventLoop) -> None:
asyncio.set_event_loop(loop)
loop.run_forever()
def run(self, coro: Coroutine[Any, Any, None]) -> None:
future = asyncio.run_coroutine_threadsafe(coro, self.get_loop())
future.result()
celery_event_loop = CeleryEventLoop()
from celery import Celery
celery_app = Celery(
broker="REDIS_URL",
include=["celery_tasks"],
)
import asyncio
from fastapi import APIRouter, FastAPI()
from celery_tasks import get_user_task, monitor_threads_task
router = APIRouter()
@router.post("/celery_tasks")
async def run_celery_tasks():
monitor_threads_task.delay()
for i in range(20):
get_user_task.delay()
await asyncio.sleep(0.5)
app = FastAPI()
import logging
import random
import time
import httpx
import threading
from _threads_loop import celery_event_loop
from celery_app import celery_app
logger = logging.getLogger(__name__)
httpx_client = httpx.AsyncClient()
async def get_user_data(user_id):
logger.info(f"===== старт таски ЮЗЕР")
delay = random.randint(1000, 4000)
data = await httpx_client.get(f"https://dummyjson.com/users/{user_id}?delay={delay}")
logger.info(f"===== конец таски ЮЗЕР")
@celery_app.task(name="get_user_task")
def get_user_task() -> None:
user_id = random.randint(1, 100)
celery_event_loop.run(get_user_data(user_id))
@celery_app.task(name="monitor_threads_task")
def monitor_threads_task() -> None:
for i in range(140):
logger.info(f"Активных потоков: {threading.active_count()}")
time.sleep(0.1)
celery -A celery_app:celery_app worker --pool=threads --concurrency=20 --loglevel=INFO
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment