使用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
相关产品推荐
相关产品推荐

