Business
Jobs
  • About Us
  • Solutions
    • Job Postings
      Post your job and receive qualified candidates in 48h.
    • Candidate Assessments
      500+ technical and psychological tests, plus anti-fraud.
    • Headhunting
      Tailor-made executive search from start to finish.
    • Payroll + EOR
      Payroll dispersal and EOR across 15+ LATAM countries.
  • Pricing
  • Jobs

0

179
Views
¿Usar trio nursery como generador de eventos enviados por servidor con FastAPI?

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.

Introducción al problema

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 .

que me gustaria lograr

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)

el problema real

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

over 4 years ago · Santiago Trujillo
1 answers
Answer question

0

Solución con websockets

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)

Solución con SSE

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.

over 4 years ago · Santiago Trujillo Report
Answer question
Find remote jobs

Discover the new way to find a job!

Top jobs
Top job categories
Business
Post vacancy Pricing Sales
Legal
Terms and conditions Privacy policy
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Show me some job opportunities
There's an error!