如何用Python高效处理超大数据集的字段匹配填充?
问题描述
我有一个超百万行、数GB规模且定期更新的数据集,该数据集通过将每条记录与其流动网络中的上游节点关联来建模流动特性。核心需求是:借助ID字段查找关联的上游设备,将上游设备记录中Num2列的数值写入当前行的Num2列,以此确定所处的流动网络层级。
数据样例(示意):
| DevNo | Up | Down | Num1 | Num2 |
|---|---|---|---|---|
| F1 | S1 | 1 | 1 | |
| F2 | S1 | S2 | 2 | 2 |
| S1 | F1 | F2 | ||
| S2 | F2 | S4 |
当前使用的代码:
import pandas as pd import numpy as np # 读取CSV到DataFrame df = pd.read_csv("L:\\Dev_h\\xxxxx.csv") # 初始化Num2字段:Num1>0时用Num1填充 df['Num2'] = np.where(df['Num1'] > 0, df['Num1'], df['Num2']) print(df) # 筛选Num2为空的行 null_num2_rows = df[df["Num2"].isnull()] # 遍历每一行,查找上游设备的Num2值填充 for row in null_num2_rows.iterrows(): row_index, row_data = row up_field = row_data["Up"] # 匹配上游设备且Num2非空的行 matched_row = df.loc[(df["DevNo"] == up_field) & ~df["Num2"].isnull()] df.at[row_index, "Num2"] = matched_row["Num2"].iloc[0] print(df) # 保存到Excel df.to_excel("L:\\ROAD ROUTING\\xxxxx.xlsx")
该代码在小数据集上运行正常,输出示例:
DevNo,Up,Down,Num1,Num2 F1,,S1,1,1 F2,S1,S2,2,2 F3,S4,S6,3,3 F4,S8,,4,4 S1,F1,F2,,1 S2,F2,S4,,2 S3,F2,S5,,2 S4,S2,F3,,2 S5,S3,T1,,2 S6,F3,S6,,3 S7,S6,S8,,3 S8,S7,F4,,3
但处理大数据集时性能极差,会耗尽内存。尝试过Pandas分块加载,但无法解决跨块匹配的问题,求优化方案。
优化方案
1. 用字典映射替代逐行循环查找
核心思路是先把所有已知Num2的设备存入字典,批量填充空值,避免重复查询和低效循环:
import pandas as pd import numpy as np # 读取数据时指定dtype减少内存占用,字符串列用str类型,数值列按需缩小 df = pd.read_csv( "L:\\Dev_h\\xxxxx.csv", dtype={"DevNo": str, "Up": str, "Down": str, "Num1": np.int32, "Num2": np.float32} ) # 初始化Num2 df['Num2'] = np.where(df['Num1'] > 0, df['Num1'], df['Num2']) # 构建已知Num2的设备映射字典:key=DevNo,value=Num2 num2_map = df[df['Num2'].notna()].set_index('DevNo')['Num2'].to_dict() # 批量填充空值:通过Up字段匹配字典值 df['Num2'] = df.apply( lambda x: num2_map.get(x['Up'], x['Num2']) if pd.isna(x['Num2']) else x['Num2'], axis=1 ) # 处理多层级依赖(上游的Num2也是后续填充的),循环直到无空值 while df['Num2'].isna().any(): # 更新映射字典,加入新填充的Num2值 new_entries = df[df['Num2'].notna() & ~df['DevNo'].isin(num2_map)].set_index('DevNo')['Num2'].to_dict() num2_map.update(new_entries) # 再次填充空值 df['Num2'] = df.apply( lambda x: num2_map.get(x['Up'], x['Num2']) if pd.isna(x['Num2']) else x['Num2'], axis=1 ) # 大数据集优先存CSV/Parquet,避免Excel的行数和内存限制 df.to_csv("L:\\ROAD ROUTING\\xxxxx_processed.csv", index=False)
2. 分块处理解决内存不足
如果全量加载仍内存不够,先离线构建映射字典,再分块读取原数据填充,无需全量数据驻留内存:
import pandas as pd import numpy as np # 第一步:单独构建初始Num2映射字典(仅加载必要列) df_meta = pd.read_csv( "L:\\Dev_h\\xxxxx.csv", usecols=['DevNo', 'Num1', 'Num2'], dtype={"DevNo": str, "Num1": np.int32} ) df_meta['Num2'] = np.where(df_meta['Num1'] > 0, df_meta['Num1'], df_meta['Num2']) num2_map = df_meta[df_meta['Num2'].notna()].set_index('DevNo')['Num2'].to_dict() # 第二步:分块读取并处理数据 chunk_size = 100000 # 按硬件调整块大小 output_path = "L:\\ROAD ROUTING\\xxxxx_processed.csv" first_write = True for chunk in pd.read_csv( "L:\\Dev_h\\xxxxx.csv", chunksize=chunk_size, dtype={"DevNo": str, "Up": str, "Down": str, "Num1": np.int32} ): # 初始化当前块的Num2 chunk['Num2'] = np.where(chunk['Num1'] > 0, chunk['Num1'], chunk['Num2']) # 填充空值 chunk['Num2'] = chunk.apply( lambda x: num2_map.get(x['Up'], x['Num2']) if pd.isna(x['Num2']) else x['Num2'], axis=1 ) # 写入文件:第一块写表头,后续追加 chunk.to_csv(output_path, mode='w' if first_write else 'a', header=first_write, index=False) first_write = False # 若存在多层级依赖,重复上述流程几次,直到无空值
3. 广度优先遍历(BFS)处理深层级网络
如果流动网络层级较深(比如A→B→C,只有C有初始Num2),BFS遍历效率远高于循环填充:
import pandas as pd import numpy as np from collections import deque df = pd.read_csv( "L:\\Dev_h\\xxxxx.csv", dtype={"DevNo": str, "Up": str, "Down": str, "Num1": np.int32} ) df['Num2'] = np.where(df['Num1'] > 0, df['Num1'], df['Num2']) # 构建反向映射:key=上游DevNo,value=待填充的下游行索引 upstream_to_downstream = {} # 初始化队列:存入所有有初始Num2的设备 queue = deque(df[df['Num2'].notna()]['DevNo'].tolist()) # 遍历空值行,关联到对应上游节点 for idx, row in df[df['Num2'].isna()].iterrows(): up_dev = row['Up'] if up_dev not in upstream_to_downstream: upstream_to_downstream[up_dev] = [] upstream_to_downstream[up_dev].append(idx) # BFS遍历填充 while queue: current_dev = queue.popleft() current_num2 = df.loc[df['DevNo'] == current_dev, 'Num2'].iloc[0] # 处理所有依赖当前节点的下游行 if current_dev in upstream_to_downstream: for idx in upstream_to_downstream[current_dev]: df.at[idx, 'Num2'] = current_num2 # 将填充后的节点加入队列,处理它的下游 downstream_dev = df.loc[idx, 'DevNo'] if downstream_dev in upstream_to_downstream: queue.append(downstream_dev) # 处理完删除映射,避免重复操作 del upstream_to_downstream[current_dev] df.to_csv("L:\\ROAD ROUTING\\xxxxx_processed.csv", index=False)
内容的提问来源于stack exchange,提问作者Jackson Dunn
相关产品推荐
相关产品推荐

