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 的内存,无论怎样分区这种情况都不会改变。

限制:

  1. 内存的限制存在
  2. 可能会对集群其他任务的运行产生不稳定的影响

解决方案:

  1. 增加 reduce 的 JVM 内存(效果可能不好)
  2. 在 key 上面做文章:在 map 阶段将造成倾斜的 key 先分成多组,例如 aaa 这个 key,map 时随机在 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 => { 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 => { 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() } }

关键注释说明:

  1. 识别倾斜 key:在实际应用中,可以通过采样统计 key 的分布频率来识别倾斜 key。
  2. 随机后缀打散:为倾斜 key 添加随机后缀(如 aaa_1, aaa_2, aaa_3, aaa_4),将原本集中到一个 reduce 的数据分散到多个 reduce 处理。
  3. 两次聚合:第一次在打散后的 key 上进行局部聚合,第二次去除后缀恢复原始 key 进行全局聚合。
  4. 随机范围选择:随机后缀的范围(如 1-4)应根据数据倾斜程度和集群资源调整,确保每个打散后的 key 数据量均衡。
  5. 性能优化:这种方法虽然增加了一次 shuffle,但避免了单个 reduce 的内存溢出,整体执行时间更稳定。

2.2 唯一值较多

问题描述:唯一值较多,单个唯一值的记录数不会超过分配给 reduce 的内存。如果发生了偶尔的数据倾斜情况,增加 reduce 个数可以缓解偶然情况下的某些 reduce 不小心分配了多个较多记录数的情况。

解决方案:增加 reduce 个数

2.3 以上两种都无效的情况

问题描述:一个固定的组合重新定义

解决方案:自定义 partitioner

三、从业务和数据上解决数据倾斜

我们能通过设计的角度尝试解决数据倾斜问题。

3.1 有损的方法

找到异常数据,比如 IP 为 0 的数据,过滤掉

3.2 无损的方法

  1. 对分布不均匀的数据,单独计算
  2. 先对 key 做一层 hash,先将数据打散让它的并行度变大,再汇集

3.3 数据预处理

通过数据预处理来避免数据倾斜

四、平台的优化方法

工作原理与优势:

  • 避免 Shuffle:Broadcast Join 将小表广播到所有 Executor,大表数据无需移动,直接在 map 端完成 join
  • 解决数据倾斜:当 join key 分布不均时,传统 SortMergeJoin 会导致某些 reduce 任务处理大量数据,而 Broadcast Join 完全避免了 reduce 阶段
  • 性能提升:对于大小表 join,Broadcast Join 通常比 SortMergeJoin 快 10-100 倍
  • 内存考量:需要确保小表能完全放入每个 Executor 的内存,否则会引发 OOM
  1. Join 操作优化:使用 map join 在 map 端就先进行 join,免得到 reduce 时卡住

实战示例(Java Spark):

import 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万条) Dataset<Row> 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万条) Dataset<Row> 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 Dataset<Row> result1 = largeTable.join(broadcast(smallTable), "user_id"); // 方法2:通过配置自动启用 Broadcast Join(当小表小于 broadcast 阈值时) spark.conf().set("spark.sql.autoBroadcastJoinThreshold", 10485760L); // 10MB Dataset<Row> 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:默认 10MB(10485760字节),控制自动启用 Broadcast Join 的表大小阈值
    • spark.sql.broadcastTimeout:默认 300秒,广播超时时间
    • spark.sql.adaptive.enabled:启用自适应查询执行,Spark 3.0+ 可自动将 SortMergeJoin 转换为 BroadcastJoin
    • spark.sql.adaptive.localShuffleReader.enabled:启用本地 shuffle reader,减少网络传输
  • Java 版本关键代码:
    import static org.apache.spark.sql.functions.broadcast; // 显式使用 broadcast Dataset<Row> result = largeDF.join(broadcast(smallDF), "user_id"); // 或者通过配置自动优化 spark.conf().set("spark.sql.autoBroadcastJoinThreshold", 10 * 1024 * 1024L); // 10MB
  1. Group 操作优化:能先进行 group 操作的时候先进行 group 操作,把 key 先进行一次 reduce,之后再进行 count 或者 distinct count 操作
  2. 压缩优化:设置 map 端输出、中间结果压缩

五、总结与展望

5.1 数据倾斜解决思路总结

通过前文的分析,我们可以将数据倾斜的解决思路归纳为三大方向:

  1. 分而治之:这是最核心的解决思路。通过将倾斜的 key 进行拆分,让原本集中在一个任务处理的数据分散到多个任务中处理。具体方法包括:
    <ul>
  2. 随机后缀法:为倾斜 key 添加随机后缀,先分散聚合再最终合并
  3. 自定义分区:根据数据分布特点设计更合理的分区策略
  4. 增加并行度:通过增加 reduce 数量来分散处理压力
  5. 业务规避:从业务逻辑和数据源头入手,避免数据倾斜的产生:
    • 数据预处理:在数据进入计算引擎前进行清洗、过滤、采样
    • 异常数据处理:识别并处理异常值、空值、默认值等
    • 业务逻辑优化:调整业务逻辑,避免产生倾斜的数据分布
  6. 平台优化:利用计算引擎的特性来规避或缓解数据倾斜:
    • Broadcast Join:避免大表 join 时的 shuffle 操作
    • Map Join:在 map 端完成 join,避免 reduce 阶段
    • 自适应执行:利用引擎的智能优化能力

5.2 未来技术展望

随着大数据技术的发展,数据倾斜问题的处理正朝着更加智能化、自动化的方向发展:

  1. Spark AQE(自适应查询执行)
    <ul>
  2. 动态合并小分区:AQE 可以自动检测到数据倾斜,并将过小的分区合并,避免任务调度开销
  3. 动态调整 Join 策略:运行时根据数据统计信息,自动将 SortMergeJoin 转换为 BroadcastJoin
  4. 倾斜 Join 优化:Spark 3.0+ 支持自动识别倾斜的 join key,并将其拆分为多个子任务处理
  5. Flink 动态负载均衡
    • Key 分组优化:Flink 1.13+ 引入了更智能的 key 分组算法,能更好地处理倾斜数据
    • 动态重分区:根据运行时数据分布情况,动态调整数据分区策略
    • 背压感知调度:结合背压机制,自动调整任务并行度和资源分配
  6. AI 驱动的优化
    • 智能倾斜检测:利用机器学习算法预测数据分布,提前识别潜在的倾斜风险
    • 自适应参数调优:根据历史执行情况和数据特征,自动优化计算参数
    • 预测性资源分配:基于数据特征预测任务资源需求,提前分配合适资源
  7. 云原生与 Serverless 架构
    • 弹性伸缩:根据数据倾斜程度动态调整计算资源
    • 异构计算:利用 GPU、FPGA 等加速倾斜数据的处理
    • 存算分离:减少数据移动,降低网络传输带来的性能影响

5.3 实践建议

在实际工作中,建议采用以下策略应对数据倾斜:

  1. 预防为主:在数据建模和业务设计阶段就考虑数据分布问题
  2. 监控先行:建立完善的数据倾斜监控体系,及时发现并预警
  3. 分层治理:根据数据倾斜的严重程度,采用不同层级的解决方案
  4. 持续优化:随着数据量和业务变化,持续优化数据倾斜处理策略
  5. 技术选型:根据业务特点选择合适的大数据计算引擎和版本

数据倾斜是大数据处理中的经典难题,但通过合理的架构设计、业务优化和技术选型,完全可以将其影响降到最低。未来随着计算引擎的不断演进和 AI 技术的深入应用,数据倾斜问题将得到更加智能、自动化的解决。

返回列表