如何在Spark读取CSV文件时添加对应文件内行号的列?
解决方案
方案1:DataFrame API实现(推荐,全原生解析,性能最优)
该方案完全依赖Spark原生实现,不需要自定义CSV解析逻辑,性能和原生spark.read.csv完全一致,适用于Spark 2.4及以上版本。
- 第一步:调整Spark配置,保证单个CSV文件不会被拆分到多个分区,确保同文件内的行顺序和原文件完全一致
# 按实际业务中最大的CSV文件大小调整,配置值要大于单文件最大体积 spark.conf.set("spark.sql.files.maxPartitionBytes", "1024MB")
- 第二步:读取CSV文件,同时新增列标记每行所属的源文件
from pyspark.sql import SparkSession from pyspark.sql.functions import input_file_name, monotonically_increasing_id from pyspark.sql.window import Window import pyspark.sql.functions as F # CSV读取参数和原生read.csv完全一致,header、delimiter等按需调整 df = spark.read.csv( path="your/csv/path/*.csv", header=True, inferSchema=False ).withColumn("source_file", input_file_name())
- 第三步:使用窗口函数按源文件分组,生成文件内行号
# 同文件在同一个分区的前提下,monotonically_increasing_id的生成顺序和原文件行顺序完全一致 window_spec = Window.partitionBy("source_file").orderBy("m_id") df_with_row_num = df.withColumn("m_id", monotonically_increasing_id())\ .withColumn("file_row_num", F.row_number().over(window_spec))\ .drop("m_id")
如果你需要的行号包含表头行,读取时设置
header=False即可,生成的行号1会对应原文件的表头行,按需过滤即可。
该方案分布式完全安全,只要保证单文件不拆分,同文件行顺序就稳定,生成的行号和原文件完全匹配。
方案2:RDD zipWithIndex实现(兼容低版本Spark)
该方案解决了zipWithIndex要求分区排序稳定的问题,通过单文件独立处理的逻辑保证行号准确,适用于Spark 2.0及以上所有版本。
- 第一步:用
wholeTextFiles读取所有CSV文件,返回「文件路径,文件全部内容」的RDD
csv_rdd = spark.sparkContext.wholeTextFiles("your/csv/path/*.csv")
- 第二步:对每个文件单独拆分、生成行号
def process_single_file(file_data): file_path, content = file_data lines = content.splitlines() # 有表头则跳过第一行,行号从1开始对应第一行数据;不需要跳表头则直接遍历lines即可 data_lines = lines[1:] return [(line, file_path, idx+1) for idx, line in enumerate(data_lines)] processed_rdd = csv_rdd.flatMap(process_single_file)
- 第三步:用原生CSV解析能力转为DataFrame
# 按实际业务定义CSV的schema schema = "col1 string, col2 int, col3 timestamp" # from_csv为Spark原生解析函数,性能和read.csv完全一致 df = processed_rdd.toDF(["raw_line", "source_file", "file_row_num"])\ .select(F.from_csv("raw_line", schema).alias("data"), "source_file", "file_row_num")\ .select("data.*", "source_file", "file_row_num")
内容的提问来源于stack exchange,提问作者vanhooser
相关产品推荐
相关产品推荐

