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

使用PySpark高效读取CSV去重行的可行性与优化方案问询

背景

作为PySpark与大数据新手,我把1300万行的CSV读入pandas DataFrame后,发现有10%的重复行,用pd.drop_duplicates()没法完全去除。试过拆分数据再去重拼接的方式,但效率极低,而且后续要处理大量同规模的CSV,所以想通过PySpark的并行化能力实现更高效的解决方案。

问题
  1. 读取1GB以上大CSV文件并去除重复行的高效方案是什么?
  2. 能不能用PySpark直接读取CSV中的去重行?这种方式比先读全量数据再调用PySpark的distinct()后写入新文件更高效吗?我直觉是并行读取去重行能降低迭代时的内存占用,但不确定是不是有认知误区。
代码示例与问题

--更新-- 我已经装好PySpark及依赖,但运行测试代码时出现ConnectionRefusedError报错,还收到任务大小超出推荐值的警告,怀疑没做并行化处理。

以下是我理解的“先读全量再去重”的实现代码:

df = pd.read_csv("<data>.csv", sep=",")

# 创建Spark Session
spark = SparkSession.builder.appName('sparkdf').getOrCreate()

# 转为Spark DataFrame
df_spark = spark.createDataFrame(df)

# 获取去重后的数据
df_distinct = df_spark.distinct()

# 将去重后的数据写入新CSV
df_distinct.to_csv("<filename>.csv", index=False)

我试过拆分数据后调用drop_duplicates()再拼接,也研究过spark.read.format这类PySpark方法。


解决方案

一、高效读取大CSV并去重的正确姿势

你当前代码的核心问题是先用pandas读全量数据再转Spark DataFrame,完全浪费了Spark的分布式能力——pandas是单机内存处理,1300万行的CSV很容易撑爆单机内存,还会导致后续Spark处理时没有并行度,进而出现任务过大警告和连接错误。

正确流程是全程用PySpark分布式处理:

from pyspark.sql import SparkSession

# 初始化Spark Session,根据资源配置并行度参数
spark = SparkSession.builder \
    .appName("CSVDeDuplication") \
    .master("local[*]")  # 本地模式下用所有可用核心;集群环境可去掉此配置
    .config("spark.executor.cores", "2")  # 每个执行器分配2核,按需调整
    .getOrCreate()

# 直接用Spark读CSV,建议手动指定Schema替代自动推断(大数据量更高效)
df_spark = spark.read \
    .option("header", "true")  # CSV带表头时开启
    .option("inferSchema", "true")  # 自动推断列类型,小数据量可用;大数据量请手动定义Schema
    .csv("<data>.csv")

# 分布式去重,Spark会将任务拆分到多个节点执行
df_distinct = df_spark.distinct()

# 写入去重后的CSV,Spark默认生成多分区文件
df_distinct.write \
    .option("header", "true") \
    .mode("overwrite")  # 覆盖已有输出目录
    .csv("<output_directory>")

二、关于“直接读取去重行”的疑问

不存在“直接读取去重行”的方法——任何工具都必须先读取数据才能判断重复。你的直觉部分正确:用Spark分布式读取+去重,确实比pandas单机处理内存占用低,因为Spark会把数据拆分到多个节点/执行器,每个节点只处理部分数据,不会把全量数据加载到单机内存。

对比两种流程:

  • 你的原代码:pandas读全量→转Spark→去重→写入。此流程中pandas已把所有数据加载到单机内存,不仅容易OOM,转Spark时还要序列化分发全量数据,效率极低。
  • 正确流程:Spark分布式读CSV→分布式去重→分布式写入。全程数据分散在多个节点,内存压力小,并行度高,效率是前者的数倍甚至数十倍。

三、解决ConnectionRefusedError和任务过大警告

  1. ConnectionRefusedError:本地模式下,确保Spark安装正确,初始化Session时加上.master("local[*]")即可;集群模式下检查集群地址配置是否正确,集群服务是否正常启动。
  2. 任务大小超出推荐值:原代码把pandas的全量DataFrame转成Spark DataFrame时,数据是作为单个任务分发的,没有拆分。改用Spark直接读CSV,Spark会自动按文件大小拆分数据为多个分区,每个分区对应一个任务,即可解决该警告。

四、可选优化:合并输出为单个CSV文件

Spark默认会生成多个分区文件,若需要单个文件,可先repartition(1)再写入,但注意这会把所有数据集中到一个节点,大数据量时可能有内存压力,仅在必要时使用:

df_distinct.repartition(1) \
    .write \
    .option("header", "true") \
    .mode("overwrite") \
    .csv("<single_file_output_dir>")

之后手动将目录里的part-*.csv重命名为目标文件名即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 03:35:41