使用PySpark高效读取CSV去重行的可行性与优化方案问询
作为PySpark与大数据新手,我把1300万行的CSV读入pandas DataFrame后,发现有10%的重复行,用pd.drop_duplicates()没法完全去除。试过拆分数据再去重拼接的方式,但效率极低,而且后续要处理大量同规模的CSV,所以想通过PySpark的并行化能力实现更高效的解决方案。
- 读取1GB以上大CSV文件并去除重复行的高效方案是什么?
- 能不能用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和任务过大警告
- ConnectionRefusedError:本地模式下,确保Spark安装正确,初始化Session时加上
.master("local[*]")即可;集群模式下检查集群地址配置是否正确,集群服务是否正常启动。 - 任务大小超出推荐值:原代码把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

