Tengo un servicio web asíncrono (con FastAPI) que ejecuta una tarea de sincronización en segundo plano.
La idea de comenzar el código es:
El código asíncrono de FastAPI recibe una solicitud de tarea.
Esta solicitud tiene un indicador de evento
Agrega la solicitud a una cola de prioridad y luego espera hasta que se establece el evento de la tarea o se alcanza un tiempo de espera.
Un trabajador de hilo de sincronización en segundo plano finaliza todas las tareas en la cola y establece el evento.
¿Cómo puedo pasar el evento de manera que el código asincrónico pueda esperarlo de manera confiable?
Ejemplo de código:
import asyncio from fastapi import FastAPI import time import uvicorn from random import randint from threading import Thread import asyncio from asyncio.locks import Event from queue import PriorityQueue, Empty global_pq = PriorityQueue() global_task_count = 0 def background_thread(): while True: time.sleep(0.5) try: task: MyTask task = global_pq.get(block=None) print(f'background {time.time()} {task}') task.complete() except Empty: print(f'background {time.time()} {None}') bg = Thread(target=background_thread) bg.setDaemon(True) bg.start() class MyTask: def __init__(self, no): self.event = Event() self.priority = randint(1, 5) self.no = no def __lt__(self, other): return self.priority < other.priority def __repr__(self): return f'Task No {self.no}, Priority {self.priority}' def complete(self): self.event._loop.call_soon_threadsafe(self.event.set) # self.event.set(True) app = FastAPI() @app.get("/") async def root(): global global_task_count global_task_count += 1 mytask = MyTask(no=global_task_count) global_pq.put((mytask.priority, mytask)) with asyncio.wait_for(mytask.event.wait(), timeout=5.0): pass return {"count": global_task_count} if __name__ == '__main__': uvicorn.run(app, host='0.0.0.0', port=9000)