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

使用PySpark获取df1中未在df2存在的行及对应行号

PySpark 大CSV差集(含行号)实现方案

针对GB级无表头单列CSV的差集计算需求,全程使用Spark分布式算子实现,避免Driver端OOM,性能最优。

核心思路

  • 先为df1生成行号标识
  • 使用Spark原生的left_anti连接实现差集计算,该算子是Spark专门为"取左表中不存在于右表的行"场景优化的,比subtract、left join后过滤null的方案效率更高,且不会错误去重左表的重复行。

代码实现

版本1:高性能版(推荐)

使用monotonically_increasing_id()生成全局唯一递增ID作为行标识,无shuffle开销,性能最好,适合不需要严格对齐原始文件物理行号、只需要唯一标识每行的场景:

from pyspark.sql import SparkSession
from pyspark.sql.functions import monotonically_increasing_id

spark = SparkSession.builder.getOrCreate()

file1 = 'file_path'
file2 = 'file_path'

# 读取无表头CSV,默认内容列列名为_c0
df1 = spark.read.csv(file1)
df2 = spark.read.csv(file2)

# 为df1添加全局唯一递增行ID
df1_with_rowid = df1.withColumn("row_num", monotonically_increasing_id())

# 左反连接取差集,结果包含行号和对应内容
diff_result = df1_with_rowid.join(
    df2,
    on="_c0",
    how="left_anti"
).select("row_num", "_c0")

# 预览结果
diff_result.show(truncate=False)

# 结果落盘(分布式写入,支持TB级数据)
# diff_result.write.csv("diff_output_path", header=False)

版本2:严格原始行号版

如果需要行号和原始CSV文件从1开始的物理行号严格对齐,使用RDD的zipWithIndex实现,仅需要一次全量数据扫描,GB级数据下开销可接受:

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, StringType, LongType

spark = SparkSession.builder.getOrCreate()

file1 = 'file_path'
file2 = 'file_path'

df1 = spark.read.csv(file1)
df2 = spark.read.csv(file2)

# 生成从1开始的连续行号
rdd_with_index = df1.rdd.zipWithIndex().map(lambda item: (item[0][0], item[1] + 1))
# 转换回DataFrame
custom_schema = StructType([
    StructField("_c0", StringType(), True),
    StructField("row_num", LongType(), False)
])
df1_with_rowid = spark.createDataFrame(rdd_with_index, schema=custom_schema)

# 取差集逻辑和版本1一致
diff_result = df1_with_rowid.join(df2, on="_c0", how="left_anti").select("row_num", "_c0")
diff_result.show(truncate=False)

注意事项

  • 不要使用df1.subtract(df2)实现差集:该算子会对结果去重,如果df1存在重复内容行,会被错误丢弃,left_anti连接会保留左表所有符合条件的行,和原始数据一致
  • 如果内容列存在null值,join时默认null不会和任何值匹配(包括null),如果需要把null也纳入匹配逻辑,可以提前用fillna给两个df的_c0列填充统一的占位值
  • GB级数据量下不要调用collect()、toPandas()把数据拉到Driver本地,会直接触发OOM,所有计算和写入都用DataFrame原生分布式API即可

内容的提问来源于stack exchange,提问作者Saad Hasan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 06:51:38