You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.26 09:54:06