You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flask中使用asyncio Futures处理异步任务的异常问题求助

Flask + Asyncio 优先级队列嵌入API问题解决

问题背景

我写了一个Flask API,用来接收句子并通过SentenceTransformer生成嵌入向量。核心逻辑是把任务加入asyncio优先级队列,高优先级任务优先处理,低优先级任务在队列长度超过10时拒绝,用Futures异步返回结果。但Flask和asyncio兼容性差,遇到两个核心问题:

  1. 用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

  1. 用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 06:05:58