Flask-RESTful中基于异步请求的分块响应发送问题修复
我搭建了一个Flask-RESTful服务,需要向前端分块发送数据以展示服务端的处理进度,进度依赖异步请求的执行结果。原本想直接用异步生成器或者类似Express.js的响应对象,但没找到可行方案,于是尝试把异步生成器转为同步。
当前逻辑是:在do函数里发起异步请求并添加到Request_Tracker(调用其add方法),请求成功后结果会被加入Request_Tracker的responses,同时调用Wait_For_Update的update方法。Wait_For_Update内部维护一个任务,update方法会取消当前任务并创建新任务。
现在遇到的问题是:generate_sync会等所有请求完成后才尝试从Request_Tracker获取下一个结果。我推测是事件循环被请求占用,run_until_complete会等循环内所有任务完成才继续执行,请问该怎么解决?
代码架构
route.py
class Class_Name(Resource): async def post_async(self, loop): requests_data = { "total": float, "request_tracker": Request_Tracker(loop) } async def do(): #...执行异步请求操作,部分请求会同时进行 await requests_data["request_tracker"].finish() async def generate(): while not requests_data["request_tracker"].get_is_finished(): yield await requests_data["request_tracker"].await_next() loop.create_task(do()) def generate_sync(): gen = generate() while True: try: yield loop.run_until_complete(gen.__anext__()) except StopAsyncIteration as e: break return generate_sync() def post(self): loop = asyncio.new_event_loop() response_generator = loop.run_until_complete(self.post_async(loop)) return app.response_class(stream_with_context(response_generator), mimetype="txt") api.add_resource(Class_Name, '/endpoint')
request_tracker.py
async def wait_forever(): await asyncio.sleep(1000000) class Wait_For_Update: def __init__(self, loop): self.loop = loop self._new() self.is_updating = False def _new(self): self.task = self.loop.create_task(wait_forever()) async def cancel(self): self.task.cancel() await asyncio.sleep(0) async def update(self): # 检查is_updating避免重复取消任务 if self.is_updating: return self.is_updating = True await self.cancel() self._new() self.is_updating = False class Request_Tracker: def __init__(self, loop): self.is_finished = False self.responses = Responses() self.wait_for_update = Wait_For_Update(loop) self.loop = loop async def add(self, req): res = await req self.responses.add(res) await self.wait_for_update.update() return res async def finish(self): await self.wait_for_update.cancel() self.is_finished = True def get_is_finished(self): return self.is_finished and len(self.responses) == 0 async def await_next(self): while True: if self.responses.has_responses(): return self.responses if self.get_is_finished(): return await self.wait_for_update.task
解决方法
1. 改用Flask原生异步视图(推荐)
Flask 2.0+原生支持异步视图,不需要手动转同步生成器,直接用异步生成器配合stream_with_context即可,从根源避免事件循环阻塞问题。
改造后的route.py:
class Class_Name(Resource): async def post(self): loop = asyncio.get_event_loop() requests_data = { "total": float, "request_tracker": Request_Tracker(loop) } async def do(): #...执行异步请求操作,部分请求会同时进行 await requests_data["request_tracker"].finish() # 启动异步任务 loop.create_task(do()) async def generate(): while not requests_data["request_tracker"].get_is_finished(): next_response = await requests_data["request_tracker"].await_next() if next_response: yield str(next_response) # 转换为可流式传输的字符串格式 # 用stream_with_context包装异步生成器返回响应 return app.response_class(stream_with_context(generate()), mimetype="txt") api.add_resource(Class_Name, '/endpoint')
注意:需要确保依赖版本符合要求(flask>=2.0, flask-restful>=0.3.9),Flask-RESTful的Resource在Flask 2.0+环境下支持异步方法。
2. 修复同步生成器的事件循环调度
如果必须保留同步视图,需要避免在同步生成器中用run_until_complete阻塞整个事件循环,改为手动调度事件循环的单次迭代,给异步请求留出执行时间:
修改generate_sync函数:
def generate_sync(): gen = generate() while True: try: # 创建异步任务,而非直接用run_until_complete阻塞 task = loop.create_task(gen.__anext__()) # 循环检查任务状态,每次迭代让出时间片给其他异步任务 while not task.done(): loop.run_until_complete(asyncio.sleep(0.01)) yield task.result() except StopAsyncIteration: break
3. 优化Request_Tracker的等待逻辑
原Wait_For_Update的任务取消逻辑过于复杂,改用asyncio.Event可以更简洁可靠地实现等待更新的逻辑:
重构Request_Tracker:
class Request_Tracker: def __init__(self, loop): self.is_finished = False self.responses = Responses() self.update_event = asyncio.Event() # 用事件替代自定义等待类 self.loop = loop async def add(self, req): res = await req self.responses.add(res) self.update_event.set() # 有新结果时触发事件 self.update_event.clear() # 重置事件等待下一次更新 return res async def finish(self): self.is_finished = True self.update_event.set() # 唤醒所有等待的任务 def get_is_finished(self): return self.is_finished and len(self.responses) == 0 async def await_next(self): while True: if self.responses.has_responses(): return self.responses if self.get_is_finished(): return await self.update_event.wait() # 等待新结果事件
内容的提问来源于stack exchange,提问作者asharpharp

