Estoy tratando de crear un punto final de eventos enviados por el servidor con FastAPI, pero no estoy seguro de si lo que estoy tratando de lograr es posible o cómo lo haría.
Básicamente, digamos que tengo una función asíncrona run_task(limit, task) que envía una solicitud asíncrona, realiza una transacción o algo similar. Digamos que para cada tarea, run_task puede devolver algunos datos JSON.
Me gustaría ejecutar múltiples tareas (múltiples run_task(limit, task) ) de forma asíncrona, para hacerlo estoy usando trio y viveros así:
async with trio.open_nursery() as nursery: limit = trio.CapacityLimiter(10) for task in tasks: nursery.start_soon(run_task, limit, task)Y finalmente, quiero devolver los resultados de cada tarea a través de un punto final FastAPI
Al principio, simplemente creé un objeto que contenía una lista y pasé ese objeto (por referencia) a cada run_task , cuando finalizaba una tarea, empujaba los datos JSON como un diccionario y devolvía el objeto completo a través del punto final una vez que todo las tareas estaban terminadas.
Esto funciona, pero lo encuentro ineficiente, el cliente que envía la solicitud debe esperar a que finalicen todas las tareas antes de poder mostrar los datos; sin embargo, algunas tareas pueden ser bastante lentas, lo que significa que los datos obtenidos de otras tareas simplemente terminan estancados .
Cada vez que finaliza una tarea, me gustaría que la API devuelva directamente los datos de dicha tarea (que habría agregado previamente al objeto) para que el cliente pueda mostrar dichos datos en tiempo real.
Fue entonces cuando descubrí lo que eran Server-Sent Events y Web-sockets. Los eventos enviados por el servidor parecían la solución adecuada para mi problema, ya que no necesito comunicación bidireccional.
Dado que FastAPI se basa en Starlette, decidí usar sse-Starlette para crear un punto de conexión con eventos enviados por el servidor. Para hacerlo, necesito crear un punto de conexión como este.
@router.get('/stream') async def runTasks( param1: str, request: Request ): event_generator = status_event_generator(request, param1) return EventSourceResponse(event_generator) Como implica el nombre status_event_generator , sse-starlette necesita devolver un generador de eventos, y ahí es donde estoy un poco atascado. Me gustaría que el generador proporcione los datos de una tarea cuando finalice (para que el cliente pueda recibir los datos de cada tarea en tiempo real), sin embargo, las tareas están dentro de la guardería del trío asíncrono, por lo que no estoy seguro de cómo para proceder
Según ¿Funciona mal el rendimiento desde el interior de un vivero en un generador asíncrono? , parece (si entiendo correctamente) que no puedo simplemente poner un rendimiento en run_task(limit, task) y esperar que funcione
En última instancia, decidí optar por websockets en lugar de SSE, ya que me di cuenta de que necesitaba pasar un objeto como datos a mi punto final, y aunque SEE puede aceptar parámetros de consulta, tratar con objetos como parámetros de consulta era demasiado complicado.
websockets con FastAPI se basan en starlette, y son bastante fáciles de usar, implementarlos en el problema anterior se puede hacer así:
@router.websocket('/stream') async def runTasks( websocket: WebSocket ): # Initialise websocket await websocket.accept() # Receive data tasks = await websocket.receive_json() async with trio.open_nursery() as nursery: limit = trio.CapacityLimiter(10) for task in tasks: nursery.start_soon(run_task, limit, task, websocket) Para devolver datos, podemos simplemente usar await await websocket.send_json() en run_task (Este es un ejemplo simplificado, preferiblemente querrá manejar cierres de websocket y casos extremos con su vivero)
Para responder al problema original, gracias a @user3840170 y https://discuss.python.org/t/preventing-yield-inside-certain-context-managers/1091 , deberíamos poder resolver el problema abriendo un vivero en algún lugar en un ámbito más amplio que contendrá el bucle que pasa por encima del generador, y usará ese vivero en el propio generador para generar tareas en segundo plano.