Flask异步接口向asyncio Queue存数据,消费端一直阻塞无法获取
问题根源分析
你的代码核心问题在于**asyncio.Queue无法跨线程的事件循环共享使用**:
- Flask的异步路由(
async def data())运行在Flask内部管理的事件循环中 flask_queue_emit()则通过executor.submit(asyncio.run, ...)在另一个线程启动了独立的事件循环
两个事件循环完全隔离,导致两边的queue.put和queue.get操作互不感知,最终queue.get()一直阻塞。
解决方案
方案一:改用线程安全队列(简单场景首选)
直接使用Python标准库的queue.Queue(线程安全)替代asyncio.Queue,无需维护异步事件循环,适合跨线程数据传递的简单需求。
修改后的完整代码:
import queue from concurrent.futures import ThreadPoolExecutor from flask import Flask, render_template app = Flask(__name__) # 替换为线程安全的queue.Queue data_queue = queue.Queue() @app.route("/", methods=("GET", "POST")) def index(): return render_template("data.html") @app.route('/data') def data(): # 同步put,线程安全 data_queue.put("data1111") data_queue.put("data22222") return "dsjkfl" def run_flask(): # 关闭debug模式(debug模式下Flask会重启线程,可能导致队列关联失效) app.run(debug=False) def flask_queue_emit(): print("flask_queue_emit begin running") while True: # 同步get,阻塞直到有数据 package = data_queue.get() print("flask_queue_emit --- queue data:", str(package)) # 执行数据处理逻辑 data_queue.task_done() print("flask_queue_emit #########") if __name__ == '__main__': executor = ThreadPoolExecutor(max_workers=4) flask_future = executor.submit(run_flask) queue_future = executor.submit(flask_queue_emit) try: while True: pass except KeyboardInterrupt: pass finally: flask_future.cancel() queue_future.cancel() executor.shutdown(wait=False)
方案二:统一使用单个asyncio事件循环(全异步场景)
如果需要保留异步操作逻辑,可以将Flask集成到单个asyncio事件循环中,通过aiohttp的WSGI适配器实现:
步骤1:安装依赖
pip install aiohttp aiohttp-wsgi
步骤2:修改代码
import asyncio from aiohttp import web from aiohttp_wsgi import WSGIHandler from flask import Flask, render_template app = Flask(__name__) queue = asyncio.Queue() @app.route("/", methods=("GET", "POST")) def index(): return render_template("data.html") @app.route('/data') async def data(): await queue.put("data1111") await queue.put("data22222") return "dsjkfl" async def flask_queue_emit(): print("flask_queue_emit begin running") while True: package = await queue.get() print("flask_queue_emit --- queue data:", str(package)) queue.task_done() print("flask_queue_emit #########") async def main(): # 将队列处理任务加入当前事件循环 asyncio.create_task(flask_queue_emit()) # 创建WSGI处理器适配Flask应用 handler = WSGIHandler(app) # 启动aiohttp服务器 aio_app = web.Application() aio_app.add_routes([web.route('*', '/{tail:.*}', handler)]) runner = web.AppRunner(aio_app) await runner.setup() site = web.TCPSite(runner, 'localhost', 5000) await site.start() print("Server started at http://localhost:5000") # 保持事件循环运行 await asyncio.Event().wait() if __name__ == '__main__': try: asyncio.run(main()) except KeyboardInterrupt: print("Server stopped")
内容的提问来源于stack exchange,提问作者Gyso
相关产品推荐
相关产品推荐

