ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Spark Shuffle 分区原理

Spark Shuffle 分区原理 一、Shuffle 分区基础原理1. 分区控制参数Map阶段并行度总核 * 2spark.default.parallelismcore_nums * 2shuffle write /shuffle read 阶段分区数由 spark.sql.shuffle.partitions 控制Spark 2.x/3.x 默认 200大数据量场景严重不足。2. Shuffle 触发判断看 Stage 是否存在宽依赖一个父分区被多个子分区复用触发 Shufflejoin / groupBy / aggregate / distinct / repartition / 窗口函数不触发 Shufflemap / filter / select / withColumn 窄依赖物理计划出现 Exchange / ShuffleExchange 节点 → 必然触发 Shuffle3. 分区数计算经验公式推荐分区数 ≈ Shuffle 总数据量 ÷ 目标分区大小目标分区大小256MB~512MB平衡内存 OOM 与 HDFS 小文件下限不低于集群总核心数保证基础并行度上限不超过 总数据量 ÷ 64MB防止小文件爆炸4. 分区数调优利弊增大分区数优点单分区数据量变小、规避 Executor OOM、提升并行度缺点Task 和中间文件变多、加重调度与 HDFS I/O 压力二、Spark 3.x AQE配置自动自适应并行度、自动合并小分区、自动倾斜 Join 优化、根治小文件与 OOM。# 核心总开关 AQE 自适应全开 spark.sql.adaptive.enabled true; spark.sql.adaptive.logLevel WARN; # 1. 分区合并解决小Task、合并输出小文件 spark.sql.adaptive.coalescePartitions.enabled true; # 计算阶段小于64MB就合并只管计算不管最终文件基准阈值别动 spark.sql.adaptive.coalescePartitions.minPartitionSize 67108864; # 计算用足够大保证 100GB 数据计算飞快 spark.sql.shuffle.partitions 200; -- 保持默认不要改小 # 写入用AQE 自动合并小文件目标输出文件/分区理想大小 256MB spark.sql.adaptive.advisoryPartitionSizeInBytes 268435456; # AQE 合并 shuffle 分区时强制保留的「最少分区数」防止合并太狠、只剩 1~2 个大分区。 spark.sql.adaptive.coalescePartitions.minPartitionNum 1; # 【关键】写入时强制合并彻底解决小文件 spark.sql.adaptive.forceCoalesceOnInsertion.enabled true; # 开启 AQE 自适应写文件最优雅不用写 coalesce 代码 # 无需手动合并分区写入时自动把小分区合并成目标大小文件参数 spark.sql.adaptive.writeEnabledtrue spark.sql.adaptive.writeTargetFileSize128m # 2. 自动数据倾斜 Join 优化 # 自动识别并打散 Join 倾斜无需手动 salting spark.sql.adaptive.skewJoin.enabled true; # 单分区超过256MB 判定为倾斜分区 spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes 268435456; # 超过平均分区5倍 判定倾斜默认5不用改 spark.sql.adaptive.skewedPartitionFactor 5; # 3. 增强Shuffle重平衡均衡分区、固化文件大小 # 本地读 Shuffle减少网络 IO提速明显。 spark.sql.adaptive.localShuffleReader.enabled true; # Shuffle 后自动重新均衡分区大小。 spark.sql.adaptive.rebalancePartitions.enabled true; # 重平衡目标分区256MB重平衡目标分区大小和落地文件对齐不乱飘。 spark.sql.adaptive.rebalancePartitionSizeInBytes 268435456; # 4. 额外兜底防小文件、读写对齐 # 读文件按 128MB 切片和 HDFS 默认块对齐读性能最优。 spark.sql.files.maxPartitionBytes 134217728; # 控制一次任务最多能生成多少个动态分区避免动态分区产生大量小文件动态分区必加避免无限创建分区 spark.sql.dynamicPartitioning.maxNumDynamicPartitions 10000;简易配置规则数据量 50GB开启前 3 项基础 AQE 即可50GB1TBspark.sql.shuffle.partitions 设 1000~20001TB4TBspark.sql.shuffle.partitions 设 2000~4000最终落地文件统一控制在 128~256MB三、写入重分区核心说明Shuffle 聚合后数据量通常大幅压缩若不做写入重分区极易生成大量小文件拖垮 HDFS 元数据和查询性能。四、重分区三种实战方案方法一PySpark DataFrame API 方案核心要点repartition(n, 列名)触发 Shuffle数据均匀打散推荐用于写入前重分区coalesce(n)无 Shuffle仅能减少分区、数据分布不均生产谨慎使用写入动态分区必须配置 Hive 动态分区参数maxRecordsPerFile 限制单文件行数避免超大文件frompyspark.sqlimportSparkSessionfrompyspark.sql.functionsimport*frompyspark.sql.windowimportWindow# 1. 初始化SparkSession 生产标准配置sparkSparkSession.builder \.appName(Ecommerce_Order_Analysis)\.config(spark.sql.shuffle.partitions,2000)\.config(spark.sql.adaptive.enabled,true)\.config(spark.sql.adaptive.coalescePartitions.enabled,true)\.config(spark.sql.adaptive.advisoryPartitionSizeInBytes,256000000)\.config(hive.exec.dynamic.partition,true)\.config(hive.exec.dynamic.partition.mode,nonstrict)\.config(hive.exec.max.dynamic.partitions,20000)\.enableHiveSupport()\.getOrCreate()# 2. 数据源读取orders_dfspark.read.parquet(/data/orders)users_dfspark.read.parquet(/data/users)products_dfspark.read.parquet(/data/products)# 3. 业务计算 含多轮Shuffleenriched_ordersorders_df \.join(broadcast(users_df),user_id,left)\.join(products_df,product_id,left)\.groupBy(category,region,dt)\.agg(sum(amount).alias(daily_sales),countDistinct(user_id).alias(unique_customers),avg(price).alias(avg_price))# 窗口函数 触发Shufflewindow_specWindow.partitionBy(dt).orderBy(col(daily_sales).desc())enriched_ordersenriched_orders.withColumn(sales_rank,rank().over(window_spec))# 4. 写入前重分区 推荐用法final_dfenriched_orders \.repartition(300,dt,region)\.sortWithinPartitions(category)# 不推荐coalesce 数据不均匀易倾斜# final_df enriched_orders.coalesce(100)# 5. 写入Hive分区表final_df.write \.mode(overwrite)\.partitionBy(dt,region)\.option(path,/user/hive/warehouse/sales.db/daily_summary)\.option(maxRecordsPerFile,2000000)\.option(compression,snappy)\.saveAsTable(sales.daily_summary)工具函数动态计算最优分区数defget_ideal_partitions(df,target_mb256,min_p20,max_p1000):按数据量自动计算最优分区数df.cache()df_size_mbdf.rdd.map(lambdax:len(str(x))).sum()/(1024*1024)idealmax(min_p,min(max_p,int(df_size_mb/target_mb)))print(f数据量{df_size_mb:.2f}MB最优分区数{ideal})returnideal# 调用示例ideal_partget_ideal_partitions(enriched_orders)final_dfenriched_orders.repartition(ideal_part,dt,region)数据倾斜手动处理工具函数defhandle_skewness(df,skew_key,total_partitions200):倾斜键加随机前缀打散适用于AQE无法自动优化场景frompyspark.sql.functionsimportconcat,lit,rand skewed_keys[category_A,category_B]df_skeweddf.filter(col(skew_key).isin(skewed_keys))df_normaldf.filter(~col(skew_key).isin(skewed_keys))df_skewed_processeddf_skewed \.withColumn(skew_key_with_prefix,concat(lit(prefix_),(rand()*10).cast(int).cast(string),lit(_),col(skew_key)))df_skewed_repartitioneddf_skewed_processed.repartition(total_partitions,skew_key_with_prefix)df_normal_repartitioneddf_normal.repartition(total_partitions,skew_key)returndf_skewed_repartitioned.union(df_normal_repartitioned)方法二Spark SQL 重分区方案核心关键字DISTRIBUTE BY按字段重分区 ShuffleDISTRIBUTE BY cast(rand() * n as int)重分区 Shuffle成n个SORT BY分区内局部排序CLUSTER BY等价 DISTRIBUTE BY SORT BY/* REPARTITION(n, 列) */Hint 强制指定分区数-- 全局参数设置SETspark.sql.shuffle.partitions2000;SETspark.sql.adaptive.enabledtrue;SEThive.exec.dynamic.partitiontrue;SEThive.exec.dynamic.partition.modenonstrict;SEThive.exec.max.dynamic.partitions10000;SEThive.exec.max.dynamic.partitions.pernode1000;-- 注册临时视图CREATEORREPLACETEMPORARYVIEWorders_viewASSELECT*FROMparquet./data/orders;CREATEORREPLACETEMPORARYVIEWusers_viewASSELECT*FROMparquet./data/users;CREATEORREPLACETEMPORARYVIEWproducts_viewASSELECT*FROMparquet./data/products;-- 业务聚合计算CREATEORREPLACETEMPORARYVIEWenriched_orders_viewASWITHjoined_dataAS(SELECTo.*,u.user_name,u.user_level,p.category,p.product_name,p.priceFROMorders_view oLEFTJOINusers_view uONo.user_idu.user_idLEFTJOINproducts_view pONo.product_idp.product_id),aggregatedAS(SELECTcategory,region,dt,SUM(amount)ASdaily_sales,COUNT(DISTINCTuser_id)ASunique_customers,AVG(price)ASavg_priceFROMjoined_dataGROUPBYcategory,region,dt)SELECT*,RANK()OVER(PARTITIONBYdtORDERBYdaily_salesDESC)ASsales_rankFROMaggregated;-- 写入重分区 方式1DISTRIBUTE BY SORT BYINSERTOVERWRITETABLEsales.daily_summaryPARTITION(dt,region)SELECTcategory,daily_sales,unique_customers,avg_price,sales_rank,dt,regionFROMenriched_orders_view DISTRIBUTEBYdt,region SORTBYcategory;-- 写入重分区 方式2Hint强制分区INSERTOVERWRITETABLEsales.daily_summaryPARTITION(dt,region)SELECT/* REPARTITION(300, dt, region) */category,daily_sales,unique_customers,avg_price,sales_rank,dt,regionFROMenriched_orders_view;-- 写入重分区 方式3CLUSTER BY 等价分区排序INSERTOVERWRITETABLEsales.daily_summaryPARTITION(dt,region)SELECTcategory,daily_sales,unique_customers,avg_price,sales_rank,dt,regionFROMenriched_orders_view DISTRIBUTEBYdt,region CLUSTERBYcategory;分桶表创建与写入-- 创建分桶表CREATETABLEIFNOTEXISTSsales.bucketed_summary(category STRING,daily_salesDOUBLE,unique_customersBIGINT,avg_priceDOUBLE,sales_rankINT)PARTITIONEDBY(dt STRING,region STRING)CLUSTEREDBY(category)INTO50BUCKETS STOREDASORC TBLPROPERTIES(orc.compressSNAPPY,transactionaltrue);-- 写入分桶表INSERTOVERWRITETABLEsales.bucketed_summaryPARTITION(dt,region)SELECTcategory,daily_sales,unique_customers,avg_price,sales_rank,dt,regionFROMenriched_orders_view;方法三动态分区 分桶写入优化# 动态分区多级分区 重分区控制小文件transformedDF \.repartition(200,col(dt),col(category)).write.mode(overwrite).partitionBy(dt,category).option(maxRecordsPerFile,1000000).option(compression,snappy).saveAsTable(target_table)# 分桶表写入 适合高频关联查询场景transformedDF \.write.mode(overwrite).bucketBy(50,category).sortBy(category).saveAsTable(bucketed_table)五、监控与验证 SQL / 代码查看执行计划final_df.explain(True)表统计信息收集ANALYZETABLEsales.daily_summaryCOMPUTESTATISTICS;核查各分区文件数量SELECTdt,region,COUNT(DISTINCTINPUT__FILE__NAME())asfile_countFROMsales.daily_summaryGROUPBYdt,regionORDERBYfile_countDESC;数据倾斜分区检测SELECTdt,region,COUNT(*)asrow_count,AVG(daily_sales)asavg_salesFROMsales.daily_summaryGROUPBYdt,regionHAVINGCOUNT(*)1000000ORDERBYrow_countDESCLIMIT10;广播 Join 优化 HintSELECT/* MAPJOIN(small_table) */*FROMlarge_tableJOINsmall_tableONlarge_table.keysmall_table.key;六、注意事项禁止滥用 coalesce无 Shuffle 合并数据分布不均极易引发数据倾斜写入前务必按Hive 分区字段做 repartition是根治小文件最有效手段必须配置 maxRecordsPerFile避免单文件过大影响查询效率Spark3.x 强制开启 AQE自动合并分区、自动倾斜优化减少人工调参动态分区作业必须配置 Hive 动态分区参数否则报错或分区生成异常小文件标准治理方案AQE开启 写入前repartition分区列 合理文件大小控制小表优先使用 broadcast 广播 Join直接消除 Shuffle 开销分桶表适合大表高频 Join、聚合场景可大幅减少查询 Shuffle避免过度分区太多小文件会影响 HDFS 和查询性能考虑存储格式ORC/Parquet 有最小文件大小要求监控资源使用观察 Executor 内存和 GC 情况结合业务特点根据数据倾斜情况调整分区策略七、调优口诀Shuffle 并行度50GB 内 10001TB 内 20004TB 左右 4000AQE 三件套全开自适应 分区合并 倾斜优化写入必做repartition 按分区字段重分区文件大小严控统一 128~256MB小文件根源未做写入前重分区 AQE 未开启
返回列表