使用Dask多核处理时Jupyter Notebook进度条不显示的问题
解决Dask在Jupyter Notebook中进度条不显示的问题
问题根源
- 原代码中
b.compute()会一次性触发所有任务并等待全部完成,后续tqdm遍历的是已完成的结果列表,无法追踪实时任务进度。 - Dask Diagnostics的
ProgressBar与tqdm同时使用存在冲突,导致进度条无法正常渲染。 - 任务提交方式未充分利用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
相关产品推荐
相关产品推荐

