使用DASK重塑多列不同的多文件时遇列不匹配错误求解决
使用DASK处理多列结构不同的Excel文件并重塑数据
问题描述
需要用DASK处理多份列结构不同的Excel文件(不同月份的文件,日期列存在差异),完成数据重塑,但运行代码时在categorize()环节触发compute()操作时出现列不匹配错误。
错误信息
ValueError: The columns in the computed data do not match the columns in the provided metadata Extra:['01/02/2022.part1', '01/02/2022.part2', '02/02/2022.part1', '02/02/2022.part2', ..., '31/02/2022.part1', '31/02/2022.part2'] Missing:['01/01/2022.part1', '01/01/2022.part2', '02/01/2022.part1', '02/01/2022.part2', ..., '31/01/2022.part1', '31/01/2022.part2']
文件结构
file1(2022年1月文件)
| Unnamed: 0 | Unnamed: 1 | 01/01/2022.part1 | 01/01/2022.part2 | 02/01/2022.part1 | 02/01/2022.part2 | ... | 31/01/2022.part1 | 31/01/2022.part2 |
|---|---|---|---|---|---|---|---|---|
| product_name | product_code | money | quantity | money | quantity | ... | money | quantity |
| product_name | product_code | money | quantity | money | quantity | ... | money | quantity |
file2(2022年2月文件)
| Unnamed: 0 | Unnamed: 1 | 01/02/2022.part1 | 01/02/2022.part2 | 02/02/2022.part1 | 02/02/2022.part2 | ... | 28/02/2022.part1 | 28/02/2022.part2 |
|---|---|---|---|---|---|---|---|---|
| product_name | product_code | money | quantity | money | quantity | ... | money | quantity |
| product_name | product_code | money | quantity | money | quantity | ... | money | quantity |
期望输出
| product_name | money | quantity |
|---|---|---|
| ball,001,01/01/2022 | 30,000 | 10,000 |
| ball,001,01/01/2022 | 15,000 | 10,000 |
现有代码
def convert_1d_to_2d(l, cols): return [l[i:i + cols] for i in range(0, len(l), cols)] def read_excel(inputs, **kwargs): return from_map(pd.read_excel, inputs, **kwargs) files = glob.glob(r'/content/*.xlsx') ddf = read_excel(files) # Columns required for pivot_table: i:columns, product_name:index ddf['i'] = str(1) ddf['product_name'] = np.nan ddf=ddf.categorize(columns=['product_name','i']) product_columns=ddf.columns[:2] date_list = ddf.columns[2:-2] date_list = convert_1d_to_2d(date_list, 11) # I want to pivot_table so I turn it by date for i in range(len(date_list)): columns_lis = [] columns_list.append(list(product_columns)) columns_list.append(list(date_list[i])) columns_list = list(itertools.chain.from_iterable(columns_list)) # Combine columns because only one index column can be specified for dask.reshape.pivot_table ddf['product_name'] = ddf['Unnamed: 0']+','+ddf['Unnamed: 1']+','+columns_list[2] dff = dd.reshape.pivot_table(ddf, values=columns_list[2:], index='product_name', columns='i')
解决方案
错误根源是不同文件列结构不一致,DASK默认用第一个文件的列作为元数据,后续文件列不匹配就会触发报错。正确思路是先对单个文件做宽表转长表处理,统一输出结构后再用DASK合并。
修正后的代码
import dask.dataframe as dd import pandas as pd import glob import re def process_single_excel(file_path): # 读取单个Excel文件 df = pd.read_excel(file_path) # 重命名产品相关列 product_df = df[['Unnamed: 0', 'Unnamed: 1']].rename(columns={ 'Unnamed: 0': 'product_name', 'Unnamed: 1': 'product_code' }) # 筛选日期相关列 date_cols = [col for col in df.columns if col not in ['Unnamed: 0', 'Unnamed: 1']] # 按日期分组(每两列对应一个日期的money和quantity) date_groups = [date_cols[i:i+2] for i in range(0, len(date_cols), 2)] processed_dfs = [] for cols in date_groups: # 从列名中提取日期(去除.part1/.part2后缀) date = re.sub(r'\.part\d+$', '', cols[0]) # 重命名money和quantity列 temp_df = df[cols].rename(columns={ cols[0]: 'money', cols[1]: 'quantity' }) # 合并产品信息、日期信息和数值列 temp_df = pd.concat([product_df, temp_df], axis=1) temp_df['full_product_id'] = temp_df['product_name'] + ',' + temp_df['product_code'] + ',' + date processed_dfs.append(temp_df) # 合并当前文件所有日期的处理结果,返回指定列 return pd.concat(processed_dfs, ignore_index=True)[['full_product_id', 'money', 'quantity']] # 获取所有Excel文件路径 files = glob.glob(r'/content/*.xlsx') # 用DASK并行处理所有文件并合并结果 ddf = dd.from_map(process_single_excel, files) # 触发计算并查看结果 result = ddf.compute() print(result.head())
代码说明
- 单文件处理:每个文件单独读取后,先提取固定的产品列,再将日期相关的宽列按每两列一组转成长表,同时从列名中提取日期,生成符合期望格式的
full_product_id列。 - 统一结构:每个文件处理后输出的列完全一致,DASK可以正确识别元数据,避免列不匹配错误。
- 并行合并:通过
dd.from_map实现多文件并行处理,自动合并所有结果,保证处理效率。
内容的提问来源于stack exchange,提问作者kaai
相关产品推荐
相关产品推荐

