Flask中使用asyncio Futures处理异步任务的异常问题求助
问题背景
我写了一个Flask API,用来接收句子并通过SentenceTransformer生成嵌入向量。核心逻辑是把任务加入asyncio优先级队列,高优先级任务优先处理,低优先级任务在队列长度超过10时拒绝,用Futures异步返回结果。但Flask和asyncio兼容性差,遇到两个核心问题:
- 用
await替代loop.run_until_complete时,报错:
RuntimeError: Task <Task pending name='Task-9' coro=<AsyncToSync.main_wrap() running at /home/user/miniconda3/envs/tf/lib/python3.9/site-packages/asgiref/sync.py:353> cb=[_run_until_complete_cb() at /home/user/miniconda3/envs/tf/lib/python3.9/asyncio/base_events.py:184]> got Future
attached to a different loop
- 用
loop.run_until_complete时,频繁报错:
Traceback (most recent call last):
File "/home/user/miniconda3/envs/tf/lib/python3.9/site-packages/flask/app.py", line 2529, in wsgi_app
response = self.full_dispatch_request()
File "/home/user/miniconda3/envs/tf/lib/python3.9/site-packages/flask/app.py", line 1825, in full_dispatch_request
rv = self.handle_user_exception(e)
File "/home/user/miniconda3/envs/tf/lib/python3.9/site-packages/flask/app.py", line 1823, in full_dispatch_request
rv = self.dispatch_request()
File "/home/user/miniconda3/envs/tf/lib/python3.9/site-packages/flask/app.py", line 1799, in dispatch_request
return self.ensure_sync(self.view_functions[rule.endpoint])(**view_args)
File "/home/user/projects/vector-crawler/backend/src/embed.py", line 78, in embed_sentences
result = loop.run_until_complete(add_to_prio_queue(sentences, prio))
File "/home/user/miniconda3/envs/tf/lib/python3.9/asyncio/base_events.py", line 623, in run_until_complete
self._check_running()
File "/home/user/miniconda3/envs/tf/lib/python3.9/asyncio/base_events.py", line 583, in _check_running
raise RuntimeError('This event loop is already running')
RuntimeError: This event loop is already running
另外程序退出时会崩溃(优先级低),求优雅处理Flask中异步任务的方案。
原代码
from flask import Flask, request, jsonify from sentence_transformers import SentenceTransformer import asyncio loop = asyncio.get_event_loop() app = Flask(__name__) prio_queue = asyncio.PriorityQueue() # Load the SentenceTransformer model model = SentenceTransformer('all-MiniLM-L6-v2') async def embed_task_loop(): while True: # get next item prio, p= await prio_queue.get() sentences, fut = p try: # Encode the sentences using the model embeddings = model.encode(sentences) # Create a dictionary to hold the sentence-embedding pairs result = { 'texts': sentences, 'embeddings': embeddings.tolist() } #return jsonify(result), 200 fut.set_result((jsonify(result), 200)) except Exception as e: #return jsonify(error=str(e)), 500 fut.set_result((jsonify(error=str(e)), 500)) async def add_to_prio_queue(sentences, prio): global prio_queue # add to prio queue always if prio is one, if prio is zero add only if prio queue is not larger than 10 if prio == 1 or (prio == 0 and prio_queue.qsize() < 10): fut=loop.create_future() package = (prio, (sentences, fut)) prio_queue.put_nowait(package) else: fut.set_result((jsonify(error='Too many requests'), 429)) return await fut @app.route('/embed', methods=['POST']) def embed_sentences(): # Get the list of sentences from the request body data = request.get_json(force=True) sentences = data.get('texts', []) prio= data.get('prio', 0) if not sentences: return jsonify(error='No sentences provided'), 400 result = loop.run_until_complete(add_to_prio_queue(sentences, prio)) return result if __name__ == '__main__': # start the embed task loop with app.app_context(): loop.create_task(embed_task_loop()) app.run()
解决办法
方案一:改用Flask异步视图 + 统一事件循环
Flask 2.0+原生支持异步视图,无需手动调用run_until_complete,只需确保所有异步操作绑定到当前运行的事件循环即可。
修改后的完整代码:
from flask import Flask, request, jsonify from sentence_transformers import SentenceTransformer import asyncio app = Flask(__name__) prio_queue = asyncio.PriorityQueue() model = SentenceTransformer('all-MiniLM-L6-v2') async def embed_task_loop(): while True: prio, (sentences, fut) = await prio_queue.get() try: embeddings = model.encode(sentences) result = { 'texts': sentences, 'embeddings': embeddings.tolist() } fut.set_result((jsonify(result), 200)) except Exception as e: fut.set_result((jsonify(error=str(e)), 500)) finally: prio_queue.task_done() # 标记任务完成,避免队列资源泄漏 async def add_to_prio_queue(sentences, prio): if prio == 1 or (prio == 0 and prio_queue.qsize() < 10): loop = asyncio.get_running_loop() fut = loop.create_future() await prio_queue.put((prio, (sentences, fut))) return await fut else: return jsonify(error='Too many requests'), 429 @app.route('/embed', methods=['POST']) async def embed_sentences(): data = request.get_json(force=True) sentences = data.get('texts', []) prio = data.get('prio', 0) if not sentences: return jsonify(error='No sentences provided'), 400 return await add_to_prio_queue(sentences, prio) if __name__ == '__main__': # 在第一次请求前启动后台任务循环 @app.before_first_request async def start_embed_loop(): asyncio.create_task(embed_task_loop()) app.run(debug=False) # 关闭debug模式,避免多进程导致的循环冲突
方案二:同步视图 + 异步函数转换
如果必须使用同步视图,可以用asgiref.sync的async_to_sync工具自动处理事件循环适配:
修改核心部分代码:
from asgiref.sync import async_to_sync async def add_to_prio_queue(sentences, prio): if prio == 1 or (prio == 0 and prio_queue.qsize() < 10): fut = asyncio.get_running_loop().create_future() await prio_queue.put((prio, (sentences, fut))) return await fut else: return jsonify(error='Too many requests'), 429 @app.route('/embed', methods=['POST']) def embed_sentences(): data = request.get_json(force=True) sentences = data.get('texts', []) prio = data.get('prio', 0) if not sentences: return jsonify(error='No sentences provided'), 400 # 用async_to_sync包装异步函数,自动处理线程间循环匹配 return async_to_sync(add_to_prio_queue)(sentences, prio)
额外优化:解决程序退出崩溃问题
在后台任务循环中添加信号监听,避免退出时阻塞:
import signal async def embed_task_loop(): stop_event = asyncio.Event() # 注册退出信号处理 def signal_handler(_, __): stop_event.set() signal.signal(signal.SIGINT, signal_handler) signal.signal(signal.SIGTERM, signal_handler) while not stop_event.is_set(): try: # 带超时获取队列,避免退出时无限阻塞 prio, (sentences, fut) = await asyncio.wait_for(prio_queue.get(), timeout=1) except asyncio.TimeoutError: continue try: embeddings = model.encode(sentences) result = { 'texts': sentences, 'embeddings': embeddings.tolist() } fut.set_result((jsonify(result), 200)) except Exception as e: fut.set_result((jsonify(error=str(e)), 500)) finally: prio_queue.task_done()
关键说明
- Flask同步视图默认运行在多线程环境,每个线程有独立事件循环,硬绑定全局loop会导致跨循环错误
- 异步视图会自动在当前请求的事件循环中执行,从根源避免循环不匹配问题
- 关闭debug模式,因为debug模式下Flask会启动多进程,导致后台任务重复启动或循环冲突
内容的提问来源于stack exchange,提问作者user2741831

