能否用Celery为Plotly Dash仪表盘实现串行任务队列?
问题描述
我有一个基于Plotly Dash实现的TERRA卫星实时Feed仪表盘,核心代码如下:
import datetime import dash from dash import dcc, html import plotly from dash.dependencies import Input, Output from dash import ctx # 需要先安装依赖:pip install pyorbital from pyorbital.orbital import Orbital satellite = Orbital('TERRA') external_stylesheets = ['https://codepen.io/chriddyp/pen/bWLwgP.css'] app = dash.Dash(__name__, external_stylesheets=external_stylesheets) app.layout = html.Div( html.Div([ html.H4('TERRA卫星实时数据'), html.Div(id='live-update-text'), dcc.Graph(id='live-update-graph'), dcc.Interval( id='interval-component', interval=1*1000, # 毫秒级间隔,触发过于频繁 n_intervals=0 ) ]) ) # 多输入触发的回调:interval定时更新、图表relayout(修改Y轴等操作) @app.callback(Output('live-update-graph', 'figure'), Input('live-update-graph', 'relayout'), Input('interval-component', 'n_intervals')) def update_graph_live(relayout, n): if ctx.triggered_id == 'relayout': # 这里是修改Y轴的逻辑,但因为interval触发太快,这段代码始终无法执行 # * 调整Y轴的代码 * return fig else: satellite = Orbital('TERRA') data = { 'time': [], 'Latitude': [], 'Longitude': [], 'Altitude': [] } # 采集历史数据 for i in range(180): time = datetime.datetime.now() - datetime.timedelta(seconds=i*20) lon, lat, alt = satellite.get_lonlatalt(time) data['Longitude'].append(lon) data['Latitude'].append(lat) data['Altitude'].append(alt) data['time'].append(time) # 创建子图 fig = plotly.tools.make_subplots(rows=2, cols=1, vertical_spacing=0.2) fig['layout']['margin'] = {'l': 30, 'r': 10, 'b': 30, 't': 10} fig['layout']['legend'] = {'x': 0, 'y': 1, 'xanchor': 'left'} fig.append_trace({ 'x': data['time'], 'y': data['Altitude'], 'name': '海拔高度', 'mode': 'lines+markers', 'type': 'scatter' }, 1, 1) fig.append_trace({ 'x': data['Longitude'], 'y': data['Latitude'], 'text': data['time'], 'name': '经纬度轨迹', 'mode': 'lines+markers', 'type': 'scatter' }, 2, 1) return fig if __name__ == '__main__': app.run_server(debug=True)
当前遇到的问题:由于interval-component触发间隔仅1秒,过于频繁,导致用户修改Y轴的relayout回调逻辑始终无法执行。我希望实现所有回调任务加入队列、按调用顺序逐个执行的逻辑,现咨询两个问题:
- 是否可以通过Celery实现这个需求?
- 如果可行,完整的小型可运行示例是什么样的?
回答
1. Celery是否可行?
完全可以。Celery是一个分布式任务队列,能够将异步任务放入队列中,按提交顺序逐个执行,正好可以解决Dash中频繁触发的interval抢占回调资源、导致用户操作回调无法执行的问题。通过将耗时或需要排队的回调任务交给Celery处理,能保证用户触发的操作(比如修改Y轴)不会被定时任务抢占,按顺序得到执行。
2. 小型可运行示例
前置依赖
需要安装以下包:
pip install dash plotly pyorbital celery redis
这里用Redis作为Celery的消息代理,所以需要先启动Redis服务(本地安装后直接启动即可,或者用Docker容器)。
完整代码示例
1. Celery配置文件 (celery_app.py)
from celery import Celery # 初始化Celery,使用Redis作为消息代理和结果存储 celery_app = Celery( 'dash_tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/0' ) # 配置任务序列化方式 celery_app.conf.task_serializer = 'json' celery_app.conf.result_serializer = 'json' celery_app.conf.accept_content = ['json']
2. Dash主应用文件 (app.py)
import datetime import dash from dash import dcc, html, ctx, Input, Output, State import plotly from pyorbital.orbital import Orbital from celery_app import celery_app from celery.result import AsyncResult import json external_stylesheets = ['https://codepen.io/chriddyp/pen/bWLwgP.css'] app = dash.Dash(__name__, external_stylesheets=external_stylesheets) # 全局变量存储最新的图表数据和任务ID latest_fig = None current_task_id = None app.layout = html.Div( html.Div([ html.H4('TERRA卫星实时数据'), html.Div(id='task-status'), dcc.Graph(id='live-update-graph'), dcc.Interval( id='interval-component', interval=1*1000, n_intervals=0 ), # 隐藏的组件用于存储任务结果 dcc.Store(id='task-result') ]) ) # -------------------------- Celery任务定义 -------------------------- @celery_app.task def generate_satellite_figure(): """生成卫星数据图表的异步任务""" satellite = Orbital('TERRA') data = { 'time': [], 'Latitude': [], 'Longitude': [], 'Altitude': [] } for i in range(180): time = datetime.datetime.now() - datetime.timedelta(seconds=i*20) lon, lat, alt = satellite.get_lonlatalt(time) data['Longitude'].append(lon) data['Latitude'].append(lat) data['Altitude'].append(alt) data['time'].append(time.isoformat()) # 序列化时间为字符串 fig = plotly.tools.make_subplots(rows=2, cols=1, vertical_spacing=0.2) fig['layout']['margin'] = {'l': 30, 'r': 10, 'b': 30, 't': 10} fig['layout']['legend'] = {'x': 0, 'y': 1, 'xanchor': 'left'} fig.append_trace({ 'x': data['time'], 'y': data['Altitude'], 'name': '海拔高度', 'mode': 'lines+markers', 'type': 'scatter' }, 1, 1) fig.append_trace({ 'x': data['Longitude'], 'y': data['Latitude'], 'text': data['time'], 'name': '经纬度轨迹', 'mode': 'lines+markers', 'type': 'scatter' }, 2, 1) # 将fig转换为JSON可序列化格式 return json.loads(plotly.io.to_json(fig)) @celery_app.task def adjust_y_axis(fig_json, relayout_data): """调整图表Y轴的异步任务""" fig = plotly.io.from_json(json.dumps(fig_json)) # 处理relayout中的Y轴调整逻辑,比如用户拖动Y轴范围 if 'yaxis.range[0]' in relayout_data and 'yaxis.range[1]' in relayout_data: fig['layout']['yaxis']['range'] = [ relayout_data['yaxis.range[0]'], relayout_data['yaxis.range[1]'] ] # 可以扩展其他Y轴调整逻辑 return json.loads(plotly.io.to_json(fig)) # -------------------------- Dash回调 -------------------------- @app.callback( Output('task-result', 'data'), Output('task-status', 'children'), Input('interval-component', 'n_intervals'), Input('live-update-graph', 'relayout'), State('task-result', 'data'), prevent_initial_call=False ) def trigger_task(n, relayout, current_result): global current_task_id, latest_fig status_text = "" # 优先处理用户的relayout操作 if ctx.triggered_id == 'live-update-graph' and relayout is not None: if latest_fig is not None: # 提交Y轴调整任务到Celery队列 task = adjust_y_axis.delay(latest_fig, relayout) current_task_id = task.id status_text = f"正在执行Y轴调整任务,ID: {task.id}" return {'task_id': task.id}, status_text else: status_text = "无可用图表数据,无法调整Y轴" return current_result, status_text # 处理interval定时任务 elif ctx.triggered_id == 'interval-component': # 如果当前没有正在执行的任务,提交新的卫星数据任务 if current_task_id is None or AsyncResult(current_task_id).ready(): task = generate_satellite_figure.delay() current_task_id = task.id status_text = f"正在更新卫星数据,任务ID: {task.id}" return {'task_id': task.id}, status_text else: status_text = f"当前有任务在执行,等待中... 任务ID: {current_task_id}" return current_result, status_text return current_result, status_text @app.callback( Output('live-update-graph', 'figure'), Input('task-result', 'data'), prevent_initial_call=False ) def update_graph_from_task(task_data): global latest_fig if task_data and 'task_id' in task_data: result = AsyncResult(task_data['task_id']) if result.ready() and result.successful(): fig_json = result.get() latest_fig = fig_json return plotly.io.from_json(json.dumps(fig_json)) # 初始加载或任务未完成时返回最新的图表 if latest_fig is not None: return plotly.io.from_json(json.dumps(latest_fig)) # 第一次加载时生成初始图表 fig = generate_satellite_figure() latest_fig = fig return plotly.io.from_json(json.dumps(fig)) if __name__ == '__main__': app.run_server(debug=True)
运行步骤
- 启动Redis服务:
- 本地安装Redis后,执行
redis-server启动 - 或者用Docker:
docker run -d -p 6379:6379 redis
- 本地安装Redis后,执行
- 启动Celery worker:
在项目目录下执行:celery -A celery_app worker --loglevel=info - 启动Dash应用:
执行python app.py,然后访问http://localhost:8050
关键说明
- 任务优先级:回调中优先判断
relayout触发(用户操作),优先提交Y轴调整任务,保证用户操作不会被定时任务抢占。 - 状态管理:通过
dcc.Store存储任务ID,定期检查任务状态,完成后更新图表。 - 任务排队:Celery会自动将任务放入队列,按提交顺序执行,同一时间只有一个任务在worker中运行(默认配置下),避免资源冲突。
内容的提问来源于stack exchange,提问作者Cauder
相关产品推荐
相关产品推荐

