Async Dash阻塞协程放入Thread未被await及to_thread运行异常问题
Python3.9 async环境运行Dash App问题解决方案
问题背景
在async事件循环中运行属于阻塞IO的Dash App,采用asyncio.to_thread简化写法时出现运行报错。
可正常运行的原始方案
通过单独启动线程运行异步主逻辑,主线程运行Dash服务的方式实现,示例代码如下(功能为采集全球航班数据生成可视化图表):
# from flask import Flask, jsonify import asyncio from threading import Thread #Dash import dash import dash_core_components as dcc import dash_html_components as html from dash.dependencies import Input, Output import requests import plotly.graph_objects as go # ** Async Part ** async def some_print_task(): """Some async function""" while True: await asyncio.sleep(2) print("Some Task") async def another_task(): """Another async function""" while True: await asyncio.sleep(3) print("Another Task") async def async_main(): """Main async function""" await asyncio.gather(some_print_task(), another_task()) def async_main_wrapper(): """Not async Wrapper around async_main to run it as target function of Thread""" asyncio.run(async_main()) # *** Dash Part ***: app = dash.Dash() app.layout = html.Div([ # html.Div([ # html.Iframe(src="https://www.flightradar24.com", # height=500,width=200) # ]), html.Div([ html.Pre(id='counter-text',children='Active Flights Worldwide'), dcc.Graph(id='live-update-graph',style={'width':1200}), dcc.Interval( id='interval-component', interval=6000, n_intervals=0) ]) ]) counter_list = [] @app.callback( Output('counter-text','children'), [Input('interval-component','n_intervals')]) def update_layout(n): url = "https://data-live.flightradar24.com/zones/fcgi/feed.js?faa=1&mlat=1&flarm=1&adsb=1&gnd=1&air=1&vehicles=1&estimated=1&stats=1" res = requests.get(url, headers={'User-Agent' : 'Mozilla/5.0'}) data = res.json() counter = 0 for element in data["stats"]["total"]: counter += data["stats"]["total"][element] counter_list.append(counter) return "Active flights Worldwide: {}".format(counter) @app.callback( Output('live-update-graph','figure'), [Input('interval-component','n_intervals')]) def update_graph(n): fig = go.Figure(data=[ go.Scatter(x=list(range(len(counter_list))), y=counter_list, mode='lines+markers') ]) return fig if __name__ == '__main__': th = Thread(target=async_main_wrapper) th.start() app.run_server(debug=True) th.join()
报错的简化写法
尝试用asyncio.to_thread简化实现时,代码如下:
import asyncio import dash async def main(): print('In Main') async def run_dashboard(): app = dash.Dash() app.run_server('0.0.0.0', 5000, debug=False) print("Running Dash") async def run(): await asyncio.gather( asyncio.to_thread(run_dashboard), main() ) asyncio.run(run())
运行时抛出警告:
RuntimeWarning: coroutine 'run_dashboard' was never awaited handle = None # Needed to break cycles when an exception occurs.
问题原因
asyncio.to_thread仅支持传入普通同步阻塞函数,不会自动执行并await传入的协程函数。你将run_dashboard定义为async协程函数后,传入to_thread只会返回未被执行的协程对象,因此触发未await的警告,Dash服务也不会正常启动。
同时原main函数仅执行一次打印就结束,会导致整个async事件循环提前退出,程序直接终止。
修复后的代码
import asyncio import dash async def main(): # 异步业务逻辑,此处添加死循环+sleep模拟持续运行的异步任务 while True: print('In Main') await asyncio.sleep(2) # 去掉async修饰符,改为普通同步函数 def run_dashboard(): app = dash.Dash() # 此处可添加Dash的布局、回调等业务配置 app.run_server('0.0.0.0', 5000, debug=False) print("Running Dash") async def run(): await asyncio.gather( asyncio.to_thread(run_dashboard), main() ) asyncio.run(run())
内容的提问来源于stack exchange,提问作者pcrx20
相关产品推荐
相关产品推荐

