使用Dask DataFrame并行加载HDF文件报错:文件路径无法识别
问题描述
我有一个Dask DataFrame,其中digitizer_get_current_savefile列存储HDF5文件的路径。想要利用Dask单机并行能力加载这些HDF5文件,但运行代码时报错OSError: File(s) not found: a,文件路径无法被识别。但相同逻辑的Pandas版本代码可以正常运行。
DataFrame示例:
digitizer_get_current_savefile motor_get_position motor_goto_strip nevts biasvoltage z-height x-dist 0 /home/data/tct-waveforms/waveform... [-76399, -20270, -1283] 1.0 15 -200.0 -1283 -76399 1 /home/data/tct-waveforms/waveform... [-76404, -20270, -1283] 2.0 15 -200.0 -1283 -76404 2 /home/data/tct-waveforms/waveform... [-76409, -20270, -1283] 3.0 15 -200.0 -1283 -76409 3 /home/data/tct-waveforms/waveform... [-76414, -20270, -1283] 4.0 15 -200.0 -1283 -76414 4 /home/data/tct-waveforms/waveform... [-76419, -20270, -1283] 5.0 15 -200.0 -1283 -76419
报错的代码:
import dask.dataframe as dd channel = "CH0" def extract_signal(row): # Read the hdf5 data file df_data = dd.read_hdf(row["digitizer_get_current_savefile"], key=channel) # Drop all columns in the datafile that begin with "Time" (only keeping the amplitudes) df_data = df_data.loc[:, ~df_data.columns.str.startswith("Time")] return df_data.mean().max() df['signal'] = df.apply(extract_signal, axis=1)
问题原因与解决方案
错误原因:
- 在Dask DataFrame的
apply函数中使用dd.read_hdf是错误的。dd.read_hdf用于批量并行读取多个HDF5文件,而此处是处理单个文件路径,且Dask的apply中row是延迟计算对象,并非实际字符串路径,导致dd.read_hdf无法正确解析路径。 - 未指定
apply的meta参数,Dask无法推断返回值的数据类型,可能引发额外类型错误。
- 在Dask DataFrame的
修正后的代码:
import dask.dataframe as dd import pandas as pd channel = "CH0" def extract_signal(row): # 使用Pandas读取单个HDF5文件,而非Dask df_data = pd.read_hdf(row["digitizer_get_current_savefile"], key=channel) # 剔除Time开头的列 df_data = df_data.loc[:, ~df_data.columns.str.startswith("Time")] # 计算均值的最大值并返回 return df_data.mean().max() # 指定meta参数,告诉Dask返回值是float类型 df['signal'] = df.apply(extract_signal, axis=1, meta=('signal', 'float64')) # 如需立即执行计算,可调用compute() # df = df.compute()
- 额外说明:
- Dask的
apply会将函数分发到各个分区并行执行,用pd.read_hdf处理单个文件符合该场景需求。 meta参数必须指定,它定义了新增列的名称和数据类型,避免Dask进行错误的类型推断。- 如果HDF5文件体积较大,单个文件的处理可能成为瓶颈,可考虑将单个HDF5文件用Dask拆分读取,但单机场景下,用Pandas读取单个文件已足够高效。
- Dask的
内容的提问来源于stack exchange,提问作者TensaZangetsu
相关产品推荐
相关产品推荐

