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

能否用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回调逻辑始终无法执行。我希望实现所有回调任务加入队列、按调用顺序逐个执行的逻辑,现咨询两个问题:

  1. 是否可以通过Celery实现这个需求?
  2. 如果可行,完整的小型可运行示例是什么样的?

回答

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)

运行步骤

  1. 启动Redis服务:
    • 本地安装Redis后,执行redis-server启动
    • 或者用Docker:docker run -d -p 6379:6379 redis
  2. 启动Celery worker:
    在项目目录下执行:
    celery -A celery_app worker --loglevel=info
    
  3. 启动Dash应用:
    执行python app.py,然后访问http://localhost:8050

关键说明

  • 任务优先级:回调中优先判断relayout触发(用户操作),优先提交Y轴调整任务,保证用户操作不会被定时任务抢占。
  • 状态管理:通过dcc.Store存储任务ID,定期检查任务状态,完成后更新图表。
  • 任务排队:Celery会自动将任务放入队列,按提交顺序执行,同一时间只有一个任务在worker中运行(默认配置下),避免资源冲突。

内容的提问来源于stack exchange,提问作者Cauder

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 10:40:31