Spark中Distinct操作执行机制及大数据场景应对方案问询
Spark中distinct操作的执行逻辑与大数据量优化方案
你用到的去重操作示例
SQL写法:
select distinct <all columns> from table/dataframe
PySpark写法:
df.select("*").distinct()
一、distinct会把全部数据拉到单个Executor吗?
不会,distinct本质是按所有列分组的特殊group by操作,它的执行逻辑和常规group by一致,是分布式处理的:
- 每个Executor先对本地存储的数据做局部去重,过滤掉本地重复的行,减少后续需要shuffle的数据量;
- 然后根据所有列的哈希值,将数据shuffle到不同的Executor节点,每个节点只处理对应哈希分区的数据;
- 最后在每个分区内完成最终去重,合并结果。
只有在极端数据倾斜场景下(比如所有数据的哈希值恰好落到同一个分区),才会出现单个Executor负载过高的情况,这属于异常情况,不是distinct操作的默认行为。
二、数据量极大时的优化方案
如果遇到数据量过大导致Executor压力高的情况,可以从以下几个维度优化:
- 先过滤无效数据:在执行distinct之前,先用
where或filter过滤掉不需要的行(比如空值、无效记录),减少后续处理的数据规模。 - 调整shuffle分区数:设置
spark.sql.shuffle.partitions参数(默认200),根据集群CPU总核数调整为核数的2-3倍,让每个shuffle分区的数据量保持在100MB左右,避免单个分区过大压垮Executor。 - 处理数据倾斜:如果某几列的重复值过多导致哈希分区倾斜,可以用加盐法拆分负载:
from pyspark.sql import functions as F # 给倾斜列加盐,拆分分区 salted_df = df.withColumn("salt", (F.rand() * 10).cast("int")) # 局部去重 partial_distinct = salted_df.dropDuplicates() # 去掉盐值后全局去重 final_distinct = partial_distinct.drop("salt").dropDuplicates() - 利用Delta表特性:如果是Delta表,先执行
OPTIMIZE table_name ZORDER BY (col1, col2),让重复数据物理上聚集,减少shuffle时的数据传输;也可以按Delta表的分区键分批去重,降低单批次处理量。 - 分批次处理:如果数据有天然的分片键(比如时间分区、地域分区),可以循环遍历每个分片,分批执行去重后合并结果,避免一次性处理全量数据。
- 采用近似去重(业务允许时):如果不需要精确的去重结果,用
approx_count_distinct函数或者HyperLogLog算法,无需全量shuffle,性能提升显著。
内容的提问来源于stack exchange,提问作者user16798185
相关产品推荐
相关产品推荐

