尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

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

大数据中的“数据倾斜“问题分析

大数据中的“数据倾斜“问题分析 前言在大数据处理领域数据倾斜是一个常见且棘手的问题。当数据分布严重不均时少数任务会处理绝大部分数据导致整个作业执行缓慢甚至失败。本文将从数据倾斜的定义、表现、原因出发系统性地介绍从业务设计、数据预处理到平台优化等多个层面的解决方案并结合 Spark、MapReduce 等主流计算框架的实战代码帮助读者全面理解和应对数据倾斜问题。一、什么是数据倾斜数据倾斜是指数据的 key 分化严重不均造成一部分数据很多一部分数据很少的局面。1.1 数据倾斜的典型表现举个 word count 的入门例子Map 阶段形成 (aaa, 1) 的形式Reduce 阶段进行 value 相加得出 aaa 出现的次数若进行 word count 的文本有 100G其中 80G 全部是 aaa剩下 20G 是其余单词就会形成 80G 的数据量交给一个 reduce 进行相加其余 20G 根据 key 不同分散到不同 reduce这种情况就造成了数据倾斜临床反应就是 reduce 跑到 99% 然后一直在原地等着那 80G 的 reduce 跑完。1.2 数据倾斜的监控表现详细查看日志或监控界面时会发现有一个或多个 reduce 卡住各种 container 报错 OOM读写的数据量极大至少远远超过其它正常的 reduce伴随着数据倾斜会出现任务被 kill 等各种诡异的表现二、数据倾斜的原因及解决方案2.1 单个值有大量记录问题描述单个值有大量记录这种值的所有记录已经超过了分配给 reduce 的内存无论怎样分区这种情况都不会改变。限制内存的限制存在可能会对集群其他任务的运行产生不稳定的影响解决方案增加 reduce 的 JVM 内存效果可能不好在 key 上面做文章在 map 阶段将造成倾斜的 key 先分成多组例如 aaa 这个 keymap 时随机在 aaa 后面加上 1,2,3,4 这四个数字之一把 key 先分成四组先进行一次运算之后再恢复 key 进行最终运算。在 MapReduce/Spark 中该方法常用下面是一个 Spark 代码示例展示如何为倾斜 key 添加随机后缀进行打散import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ object DataSkewSolution { def main(args: Array[String]): Unit { val spark SparkSession.builder() .appName(DataSkewSolution) .master(local[*]) .getOrCreate() import spark.implicits._ // 模拟有数据倾斜的数据集 val data Seq( (aaa, 1), (aaa, 1), (aaa, 1), (aaa, 1), (aaa, 1), (aaa, 1), (aaa, 1), (aaa, 1), (aaa, 1), (aaa, 1), (bbb, 1), (ccc, 1), (ddd, 1), (eee, 1) ) val df spark.createDataset(data).toDF(key, value) // 第一步识别倾斜的 key这里假设 aaa 是倾斜 key val skewedKey aaa // 第二步为倾斜 key 添加随机后缀打散到多个分区 val dfWithSuffix df.map(row gt; { val key row.getString(0) val value row.getInt(1) if (key skewedKey) { // 为倾斜 key 添加 1-4 的随机后缀 val randomSuffix scala.util.Random.nextInt(4) 1 (s${key}_${randomSuffix}, value) } else { (key, value) } }).toDF(new_key, value) // 第三步第一次聚合在打散后的 key 上进行 val firstAgg dfWithSuffix .groupBy(new_key) .agg(sum(value).as(partial_sum)) // 第四步恢复原始 key进行最终聚合 val finalResult firstAgg.map(row gt; { val newKey row.getString(0) val partialSum row.getLong(1) if (newKey.startsWith(s${skewedKey}_)) { // 去除随机后缀恢复原始 key (skewedKey, partialSum) } else { (newKey, partialSum) } }).toDF(key, total) .groupBy(key) .agg(sum(total).as(final_count)) // 显示结果 finalResult.show() spark.stop() } }关键注释说明识别倾斜 key在实际应用中可以通过采样统计 key 的分布频率来识别倾斜 key。随机后缀打散为倾斜 key 添加随机后缀如 aaa_1, aaa_2, aaa_3, aaa_4将原本集中到一个 reduce 的数据分散到多个 reduce 处理。两次聚合第一次在打散后的 key 上进行局部聚合第二次去除后缀恢复原始 key 进行全局聚合。随机范围选择随机后缀的范围如 1-4应根据数据倾斜程度和集群资源调整确保每个打散后的 key 数据量均衡。性能优化这种方法虽然增加了一次 shuffle但避免了单个 reduce 的内存溢出整体执行时间更稳定。2.2 唯一值较多问题描述唯一值较多单个唯一值的记录数不会超过分配给 reduce 的内存。如果发生了偶尔的数据倾斜情况增加 reduce 个数可以缓解偶然情况下的某些 reduce 不小心分配了多个较多记录数的情况。解决方案增加 reduce 个数2.3 以上两种都无效的情况问题描述一个固定的组合重新定义解决方案自定义 partitioner三、从业务和数据上解决数据倾斜我们能通过设计的角度尝试解决数据倾斜问题。3.1 有损的方法找到异常数据比如 IP 为 0 的数据过滤掉3.2 无损的方法对分布不均匀的数据单独计算先对 key 做一层 hash先将数据打散让它的并行度变大再汇集3.3 数据预处理通过数据预处理来避免数据倾斜四、平台的优化方法工作原理与优势避免 ShuffleBroadcast Join 将小表广播到所有 Executor大表数据无需移动直接在 map 端完成 join解决数据倾斜当 join key 分布不均时传统 SortMergeJoin 会导致某些 reduce 任务处理大量数据而 Broadcast Join 完全避免了 reduce 阶段性能提升对于大小表 joinBroadcast Join 通常比 SortMergeJoin 快 10-100 倍内存考量需要确保小表能完全放入每个 Executor 的内存否则会引发 OOMJoin 操作优化使用 map join 在 map 端就先进行 join免得到 reduce 时卡住实战示例Java Sparkimport org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; import org.apache.spark.sql.functions; import static org.apache.spark.sql.functions.broadcast; public class BroadcastJoinExample { public static void main(String[] args) { // 创建 SparkSession SparkSession spark SparkSession.builder() .appName(BroadcastJoinExample) .master(local[*]) .getOrCreate(); // 模拟数据大表用户行为日志1000万条 DatasetRow largeTable spark.createDataFrame( Arrays.asList( RowFactory.create(1, view, 2024-01-01), RowFactory.create(2, click, 2024-01-01), RowFactory.create(3, purchase, 2024-01-01), RowFactory.create(1, click, 2024-01-02), RowFactory.create(2, purchase, 2024-01-02) // ... 更多数据 ), new StructType(new StructField[]{ DataTypes.createStructField(user_id, DataTypes.IntegerType, false), DataTypes.createStructField(action, DataTypes.StringType, false), DataTypes.createStructField(date, DataTypes.StringType, false) }) ); // 小表用户信息表1万条 DatasetRow smallTable spark.createDataFrame( Arrays.asList( RowFactory.create(1, Alice, 北京, VIP), RowFactory.create(2, Bob, 上海, 普通), RowFactory.create(3, Charlie, 广州, VIP), RowFactory.create(4, David, 深圳, 普通) // ... 更多数据 ), new StructType(new StructField[]{ DataTypes.createStructField(user_id, DataTypes.IntegerType, false), DataTypes.createStructField(name, DataTypes.StringType, false), DataTypes.createStructField(city, DataTypes.StringType, false), DataTypes.createStructField(level, DataTypes.StringType, false) }) ); // 方法1使用 broadcast hint 显式指定 Broadcast Join DatasetRow result1 largeTable.join(broadcast(smallTable), user_id); // 方法2通过配置自动启用 Broadcast Join当小表小于 broadcast 阈值时 spark.conf().set(spark.sql.autoBroadcastJoinThreshold, 10485760L); // 10MB DatasetRow result2 largeTable.join(smallTable, user_id); // 查看执行计划确认使用了 BroadcastHashJoin result1.explain(); System.out.println(结果示例); result1.show(5); spark.stop(); } }适用场景与核心参数配置Map Join / Broadcast Join 适用场景大表 join 小表小表数据量远小于大表通常小表能完全放入每个 Executor 的内存中维度表 join 事实表在数据仓库场景中维度表通常较小事实表较大解决数据倾斜当 join key 分布不均导致某些 reduce 任务过载时使用 Broadcast Join 可以避免 shuffle核心参数配置spark.sql.autoBroadcastJoinThreshold默认 10MB10485760字节控制自动启用 Broadcast Join 的表大小阈值spark.sql.broadcastTimeout默认 300秒广播超时时间spark.sql.adaptive.enabled启用自适应查询执行Spark 3.0 可自动将 SortMergeJoin 转换为 BroadcastJoinspark.sql.adaptive.localShuffleReader.enabled启用本地 shuffle reader减少网络传输Java 版本关键代码import static org.apache.spark.sql.functions.broadcast; // 显式使用 broadcast DatasetRow result largeDF.join(broadcast(smallDF), user_id); // 或者通过配置自动优化 spark.conf().set(spark.sql.autoBroadcastJoinThreshold, 10 * 1024 * 1024L); // 10MBGroup 操作优化能先进行 group 操作的时候先进行 group 操作把 key 先进行一次 reduce之后再进行 count 或者 distinct count 操作压缩优化设置 map 端输出、中间结果压缩五、总结与展望5.1 数据倾斜解决思路总结通过前文的分析我们可以将数据倾斜的解决思路归纳为三大方向分而治之这是最核心的解决思路。通过将倾斜的 key 进行拆分让原本集中在一个任务处理的数据分散到多个任务中处理。具体方法包括ul随机后缀法为倾斜 key 添加随机后缀先分散聚合再最终合并自定义分区根据数据分布特点设计更合理的分区策略增加并行度通过增加 reduce 数量来分散处理压力业务规避从业务逻辑和数据源头入手避免数据倾斜的产生数据预处理在数据进入计算引擎前进行清洗、过滤、采样异常数据处理识别并处理异常值、空值、默认值等业务逻辑优化调整业务逻辑避免产生倾斜的数据分布平台优化利用计算引擎的特性来规避或缓解数据倾斜Broadcast Join避免大表 join 时的 shuffle 操作Map Join在 map 端完成 join避免 reduce 阶段自适应执行利用引擎的智能优化能力5.2 未来技术展望随着大数据技术的发展数据倾斜问题的处理正朝着更加智能化、自动化的方向发展Spark AQE自适应查询执行ul动态合并小分区AQE 可以自动检测到数据倾斜并将过小的分区合并避免任务调度开销动态调整 Join 策略运行时根据数据统计信息自动将 SortMergeJoin 转换为 BroadcastJoin倾斜 Join 优化Spark 3.0 支持自动识别倾斜的 join key并将其拆分为多个子任务处理Flink 动态负载均衡Key 分组优化Flink 1.13 引入了更智能的 key 分组算法能更好地处理倾斜数据动态重分区根据运行时数据分布情况动态调整数据分区策略背压感知调度结合背压机制自动调整任务并行度和资源分配AI 驱动的优化智能倾斜检测利用机器学习算法预测数据分布提前识别潜在的倾斜风险自适应参数调优根据历史执行情况和数据特征自动优化计算参数预测性资源分配基于数据特征预测任务资源需求提前分配合适资源云原生与 Serverless 架构弹性伸缩根据数据倾斜程度动态调整计算资源异构计算利用 GPU、FPGA 等加速倾斜数据的处理存算分离减少数据移动降低网络传输带来的性能影响5.3 实践建议在实际工作中建议采用以下策略应对数据倾斜预防为主在数据建模和业务设计阶段就考虑数据分布问题监控先行建立完善的数据倾斜监控体系及时发现并预警分层治理根据数据倾斜的严重程度采用不同层级的解决方案持续优化随着数据量和业务变化持续优化数据倾斜处理策略技术选型根据业务特点选择合适的大数据计算引擎和版本数据倾斜是大数据处理中的经典难题但通过合理的架构设计、业务优化和技术选型完全可以将其影响降到最低。未来随着计算引擎的不断演进和 AI 技术的深入应用数据倾斜问题将得到更加智能、自动化的解决。
返回列表