PySpark实现CSV文件层级结构扁平化的方法
实现类JSON结构CSV的LIN行关联父级PER值
我们需要处理类JSON结构的CSV文件,将所有标记为LIN的行新增一列,填入其所属的最近PER行的值,LIN行按多个PER分组。
原始输入数据
"ENV","ABC" "HDR","Test message" "PER",20221210 "LIN",1 "LIN",2 "PER",20221212 "LIN",1 "LIN",2 "LIN",3 "LIN",4 "PER",20221213 "LIN",1 "LIN",2 "LIN",3 "CNT",9
期望输出结构
"ENV","ABC" "HDR","Test message" "PER",20221210 "LIN",1,20221210 "LIN",2,20221210 "PER",20221212 "LIN",1,20221212 "LIN",2,20221212 "LIN",3,20221212 "LIN",4,20221212 "PER",20221213 "LIN",1,20221213 "LIN",2,20221213 "LIN",3,20221213 "CNT",9
实现方案
方案一:使用Pandas(适合中小文件)
Pandas的向前填充(ffill)可以快速实现PER值的继承,代码简洁易读:
import pandas as pd # 读取CSV,无表头 df = pd.read_csv('input.csv', header=None) # 生成PER值列:遇到PER行记录值,其他行留空后向前填充 df['parent_per'] = df.apply(lambda x: x[1] if x[0] == 'PER' else pd.NA, axis=1).ffill() # 构造结果行:LIN行追加PER值,其他行保持原结构 result_df = df.apply( lambda row: pd.Series([row[0], row[1], row['parent_per']]) if row[0] == 'LIN' else pd.Series([row[0], row[1]]), axis=1 ) # 保存输出,确保所有字段被引号包裹 result_df.to_csv('output.csv', index=False, header=False, quoting=pd.io.common.csv.QUOTE_ALL)
方案二:纯Python逐行处理(适合超大文件)
如果文件体积过大,用纯Python逐行读写可以避免内存占用过高的问题:
input_path = 'input.csv' output_path = 'output.csv' current_per = None with open(input_path, 'r', encoding='utf-8') as infile, open(output_path, 'w', encoding='utf-8') as outfile: for line in infile: line = line.strip() if not line: continue # 拆分并去除字段引号 parts = [p.strip('"') for p in line.split(',')] if parts[0] == 'PER': current_per = parts[1] # 原样写入PER行 outfile.write(f'"{parts[0]}","{current_per}"\n') elif parts[0] == 'LIN': # 写入带PER值的LIN行 outfile.write(f'"{parts[0]}","{parts[1]}","{current_per}"\n') else: # 其他行原样拼接引号后写入 quoted_parts = [f'"{p}"' for p in parts] outfile.write(','.join(quoted_parts) + '\n')
两种方案都能实现需求:方案一适合数据量不大的场景,代码更简洁;方案二更适合处理GB级的超大CSV文件,内存效率更高。
内容的提问来源于stack exchange,提问作者Ian Smith
相关产品推荐
相关产品推荐

