使用Dask多进程构建文件摘要生成应用时无法打印输出及更新列表的问题排查
分析你的Dask多进程应用问题及修复方案
我来帮你拆解下代码里的问题,以及如何调整才能实现你想要的功能:
核心问题拆解
你的代码主要卡在Dask分布式进程的隔离特性和全局变量的误用上,具体来说:
- 全局变量跨进程无效:Dask的worker是独立的进程,主进程里的
data_date、csv_file_list这些变量,每个worker都会生成自己的副本。你在my_function里修改data_date.append(date),改的只是当前worker进程里的列表,主进程完全看不到这个变化,所以最后主进程的data_date还是空的。 - 文件名获取错误:你在函数里写
date = csv_file_list,这是把整个文件名列表赋值给date,而不是当前遍历的单个文件名,应该取csv_file_list[i]才对。 - Worker打印无法直接显示:Worker进程的标准输出不会自动同步到主进程的控制台,所以你写的
print(date)不会在你的终端里显示,得通过返回结果或者日志才能看到。 - 图表生成的限制:如果在worker里生成图表,因为worker是后台进程,没有图形界面上下文,直接调用
plt.show()会报错;而且就算生成了图表,也没法直接传递回主进程,必须把数据拿到主进程后再生成图表。
修复后的完整代码
下面是调整后的代码,我标注了关键修改点:
import pandas as pd import numpy as np import datetime as dt import matplotlib.pyplot as plt plt.ioff() import time from pathlib import Path import webbrowser from dask.distributed import Client def my_function(file_path): # 直接接收文件路径,避免依赖主进程的全局变量 df = pd.read_csv(file_path, skiprows=0) # 获取当前处理的文件名 file_name = Path(file_path).name # 这里替换成你实际的straddle价格计算逻辑 # 示例:假设数据里有open/close列,计算均值作为straddle价格 straddle_open = (df['open'] + df['close']).mean() straddle_close = (df['open'].iloc[-1] + df['close'].iloc[-1]) # 假设文件名包含日期,根据你的实际文件名格式调整解析方式 date = pd.to_datetime(file_name.split('.')[0]) # 返回处理后的结果,而非修改全局变量 return { 'Date': date, 'straddle_price_open': straddle_open, 'straddle_price_close': straddle_close } if __name__ == '__main__': client = Client(n_workers=4, threads_per_worker=2) webbrowser.open(client.dashboard_link) print(client) # 用Path管理路径,更可靠 data_dir = Path("D:\\spx\\Complete data\\item_000027392") # 获取所有zip文件的完整路径 csv_file_paths = list(data_dir.glob("*.zip")) total_file = len(csv_file_paths) print(f"发现{total_file}个待处理的zip文件") # 直接把文件路径传给worker,避免依赖全局变量 futures = [client.submit(my_function, path) for path in csv_file_paths] # 收集所有worker返回的结果 results = client.gather(futures) # 在主进程中整理结果为DataFrame summary = pd.DataFrame(results) print("\n处理完成,摘要数据预览:") print(summary.head()) # 主进程中生成图表(只有主进程有图形上下文) plt.figure(figsize=(10,6)) plt.plot(summary['Date'], summary['straddle_price_open'], label='Straddle Open') plt.plot(summary['Date'], summary['straddle_price_close'], label='Straddle Close') plt.title('Straddle Price Trend') plt.xlabel('Date') plt.ylabel('Price') plt.legend() plt.show() client.close()
关键修改说明
- 抛弃全局变量依赖:把每个文件的完整路径直接传给worker函数,让worker不需要访问主进程的变量,避免进程隔离带来的问题。
- 函数返回结果而非修改全局变量:worker处理完数据后直接返回字典,主进程通过
client.gather收集所有结果,这样就能拿到所有处理后的数据。 - 图表移到主进程生成:所有数据汇总到主进程后,再用Matplotlib绘制图表,避免worker进程的图形上下文问题。
- 添加
if __name__ == '__main__'::这是Python多进程编程的标准写法,防止子进程重复初始化Client等资源,避免意外错误。
额外小提示
- 如果想查看worker的打印信息,可以在创建Client时加上日志配置:
Client(n_workers=4, threads_per_worker=2, silence_logs=0),这样worker的日志会显示在主控制台。 - 如果你的数据集很大,建议直接用Dask DataFrame(
dd.read_csv)来处理,它会自动帮你并行读取和处理文件,比手动提交任务更高效。
内容的提问来源于stack exchange,提问作者yashaswi kashyap
相关产品推荐
相关产品推荐

