如何处理竖线分隔文件的列换行并导入PySpark DataFrame?
处理固定宽度换行的竖线分隔文件导入PySpark DataFrame
针对50GB无引号包裹、字段因固定宽度被截断换行的竖线分隔文件,不需要额外外部预处理工具,直接用PySpark就能完成合并与导入,以下是两种可行方案:
方案一:按字段分隔符数量合并行(推荐)
核心思路:每条完整记录的竖线数量固定(比如5列对应4个竖线),通过统计每行的竖线数量,将不完整的行(竖线数不足)合并到上一条记录中。
步骤代码:
- 读取原始文本行
raw_lines = spark.read.text("path/to/your/large_file.txt")
- 标记完整记录并生成分组ID
from pyspark.sql import functions as F from pyspark.sql.window import Window # 统计每行的竖线数量 lines_with_pipes = raw_lines.withColumn( "pipe_count", F.size(F.split(F.col("value"), r"\|")) - 1 ) # 标记是否为完整记录的结尾(假设完整记录含4个竖线,对应5列) lines_with_flag = lines_with_pipes.withColumn( "is_complete", F.when(F.col("pipe_count") == 4, 1).otherwise(0) ) # 生成分组ID,将属于同一条记录的行归为一组 window = Window.orderBy(F.monotonically_increasing_id()) grouped_lines = lines_with_flag.withColumn( "group_id", F.sum("is_complete").over(window.rangeBetween(Window.unboundedPreceding, -1)) ).fillna({"group_id": 0}) # 处理第一条记录的分组ID为空的情况
- 合并同组行并解析为DataFrame
# 合并同一组的行,拼接成完整记录 merged_records = grouped_lines.groupBy("group_id").agg( F.concat_ws("", F.collect_list("value")).alias("full_record") ) # 解析完整记录为结构化DataFrame from pyspark.sql.types import StructType, StructField, StringType schema = StructType([ StructField("Column1", StringType()), StructField("Column2", StringType()), StructField("Column3", StringType()), StructField("Column4", StringType()), StructField("Column5", StringType()) ]) final_df = merged_records.select( F.from_csv(F.col("full_record"), sep="|", schema=schema).alias("data") ).select("data.*")
方案二:按固定字符宽度合并行(已知换行宽度时使用)
如果明确知道每行的固定截断宽度(比如每行最多80个字符),且每条完整记录的总长度固定,可以通过累计字符长度来合并行:
步骤代码:
raw_lines = spark.read.text("path/to/your/large_file.txt") # 计算每行的字符长度 lines_with_length = raw_lines.withColumn( "line_length", F.length(F.col("value")) ) # 累计字符长度,生成分组ID(假设完整记录总长度为200) window = Window.orderBy(F.monotonically_increasing_id()) lines_with_cum_length = lines_with_length.withColumn( "cum_length", F.sum("line_length").over(window.rangeBetween(Window.unboundedPreceding, 0)) ) grouped_lines = lines_with_cum_length.withColumn( "group_id", F.floor((F.col("cum_length") - 1) / 200) # 200为完整记录总长度 ) # 合并同组行并解析,后续步骤同方案一 merged_records = grouped_lines.groupBy("group_id").agg( F.concat_ws("", F.collect_list("value")).alias("full_record") )
注意事项
- 若文件存在空行,可在读取后先过滤:
raw_lines = raw_lines.filter(F.col("value") != "") - 对于极端大文件,建议先采样小部分数据,验证完整记录的竖线数量或总长度,再应用到全量数据
- 优先使用方案一,因为它不需要预先知道固定宽度,适应性更强
内容的提问来源于stack exchange,提问作者Kalyan Tirunahari
相关产品推荐
相关产品推荐

