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

使用Dask多核处理时Jupyter Notebook进度条不显示的问题

解决Dask在Jupyter Notebook中进度条不显示的问题

问题根源

  1. 原代码中b.compute()会一次性触发所有任务并等待全部完成,后续tqdm遍历的是已完成的结果列表,无法追踪实时任务进度。
  2. Dask Diagnostics的ProgressBar与tqdm同时使用存在冲突,导致进度条无法正常渲染。
  3. 任务提交方式未充分利用Dask分布式客户端的任务追踪能力。

解决方案1:用Dask Delayed + Tqdm 追踪任务进度

改用delayed定义任务,通过客户端提交后,用tqdm逐个追踪任务完成状态,适配Jupyter的交互环境:

import pandas as pd
import os
import dask
from datetime import datetime
from dask.distributed import Client, LocalCluster
from dask import delayed
from tqdm.auto import tqdm

id_list = [1,2,3,4,5]  # 替换为你的10000个传感器ID
out_folder = 'output' + os.path.sep

def compute_parallel_dask(pid, out_folder):
    df_data = get_data(pid)
    fname = out_folder + f'{pid}.parquet'
    if not os.path.exists(fname):
         # 优化:用列表推导替代iterrows+loc,提升单任务效率
         rows = []
         for indx, row in df_data.iterrows(): 
               raw_value = row.value
               algo1_val = compute_algo1(row)
               algo2_val = compute_algo2(row)
               rows.append([pd.to_datetime(indx), raw_value, algo1_val, algo2_val])
         df_computed = pd.DataFrame(rows, columns=['DateTime','raw_value','algo1_value','algo2_value'])
         df_computed.to_parquet(fname, index=False)
    else:
        df_computed = pd.read_parquet(fname)
    return df_computed  # 直接返回pandas DataFrame,减少不必要的Dask转换开销

# 启动本地集群
num_cores = 6
cluster = LocalCluster(n_workers=num_cores, scheduler_port=0)
client = Client(cluster)

# 生成Delayed任务列表
delayed_tasks = [delayed(compute_parallel_dask)(pid, out_folder) for pid in id_list]

# 提交任务到集群并追踪进度
futures = client.compute(delayed_tasks)
results = []
for future in tqdm(futures, desc="Processing Sensors"):
    results.append(future.result())

print('All Done')
# 关闭资源
client.close()
cluster.close()

解决方案2:使用Dask原生ProgressBar

如果偏好Dask自带的进度条,需避免与tqdm冲突,用上下文管理器确保进度条在Jupyter中正常渲染:

import pandas as pd
import os
import dask
from datetime import datetime
from dask.distributed import Client, LocalCluster
from dask import bag
from dask.diagnostics import ProgressBar

id_list = [1,2,3,4,5]
out_folder = 'output' + os.path.sep

def compute_parallel_dask(pid, out_folder):
    df_data = get_data(pid)
    fname = out_folder + f'{pid}.parquet'
    if not os.path.exists(fname):
         rows = []
         for indx, row in df_data.iterrows(): 
               raw_value = row.value
               algo1_val = compute_algo1(row)
               algo2_val = compute_algo2(row)
               rows.append([pd.to_datetime(indx), raw_value, algo1_val, algo2_val])
         df_computed = pd.DataFrame(rows, columns=['DateTime','raw_value','algo1_value','algo2_value'])
         df_computed.to_parquet(fname, index=False)
    else:
        df_computed = pd.read_parquet(fname)
    return df_computed

# 启动集群
num_cores = 6
cluster = LocalCluster(n_workers=num_cores, scheduler_port=0)
client = Client(cluster)

# 创建Dask Bag并执行任务
b = bag.from_sequence(id_list).map(lambda pid: compute_parallel_dask(pid, out_folder))
with ProgressBar():
    results = b.compute()

print('All Done')
client.close()
cluster.close()

额外优化建议

  • 避免在任务函数中把pandas DataFrame转为Dask DataFrame,直接返回pandas对象可减少开销。
  • 若get_data从数据库读取数据,可尝试用Dask DataFrame直接分块读取全量数据,再按id分组处理,进一步提升分布式计算效率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 16:16:36