十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Spark核心机制与性能优化实战:从内存计算到数据倾斜排查

Spark核心机制与性能优化实战:从内存计算到数据倾斜排查 这几年大数据岗位的面试十家有九家拿 Spark 的核心机制出来压轴。生产环境里的数据开发也一样从日志清洗到实时特征工程从数仓同步到机器学习训练前的样本准备几乎处处都是 Spark 的影子。我自己做了这么多年的数据平台最大的感受是你不一定需要把 Spark 源码翻个底朝天但技术核心、应用边界、性能调优这三个维度必须打通否则线上集群就是一个随时会爆的盲盒。这篇文章是这个系列的 1.3 篇定位是综合实战向。我会重点解决三件事一是把 Spark 的内存计算和任务调度逻辑用大白话讲清楚二是拆解它在真实业务场景里到底怎么落地三是我这几年踩过的性能优化坑和排查方法整理成一套可以直接抄作业的流程。不管你是刚把 Spark 装起来的新手还是正在为 OOM 和数据倾斜头疼的老手这篇文章都会尽量把原理和操作放在一起讲保证你能看得懂、用得上。1. Spark 技术体系拆解它凭什么成为大数据计算的主流选择1.1 从 Hadoop MapReduce 到内存计算核心逻辑的演进要理解 Spark得先知道它解决了什么问题。Hadoop MapReduce 的计算模型很朴素每个作业拆成 Map 和 Reduce 两个阶段中间结果必须落到磁盘下一阶段再从磁盘读回来。这个设计在当时分布式计算刚刚起步的年代是合理的但它有一个致命的硬伤——中间结果落盘带来的 I/O 开销太大了。举个例子如果要做十轮迭代计算MapReduce 每轮都要写一次磁盘再读一次磁盘十轮下来就是二十次全量磁盘 I/O。这在 PB 级数据规模下几乎不可接受。Spark 的核心思路很直接把计算过程的中间结果尽量留在内存里用 DAG有向无环图把多个计算步骤组织起来同一份数据在内存里反复使用省掉落盘的开销。这就是它比 MapReduce 快一到两个数量级的根本原因。我见过不少初学者把 Spark 的“快”归结为“内存比磁盘快”这话对了一半。更准确地说Spark 快在两点一个是内存计算减少了 I/O 次数另一个是 DAG 调度器能把多个计算步骤在同一个 task 里串起来执行减少任务启动和调度的开销。这两点叠加在迭代计算、交互式查询、复杂 ETL 这些场景里性能差距特别明显。1.2 RDD、DataFrame、Dataset 怎么选别再做无谓的纠结这三个抽象是 Spark 面试的高频题也是实际开发里最容易犯选择困难症的地方。我的理解很简单RDD 是最底层的抽象它的优势是灵活能处理非结构化的数据能精细控制分区和缓存策略。缺点是序列化和反序列化开销大而且没有自动优化机制。现在纯粹用 RDD 写业务逻辑的场景已经很少了除非你在做非常底层的自定义算子。DataFrame 是当前最主流的选择。它带上了 schema 信息Spark 能根据 schema 做列裁剪、谓词下推这些优化。同样的逻辑DF 版本比 RDD 版本通常快 2 到 5 倍代码还更简洁。我平时做数据开发90% 的场景都是 DataFrame 加 Spark SQL。Dataset 是强类型版本主要在 Scala 里用Python 里没有对应的概念。它的价值是编译期类型检查适合复杂业务对象处理。但 Dataset 在序列化和编解码上有额外开销性能不一定比 DataFrame 好。我的建议是Java/Scala 项目里如果重视类型安全可以用 Dataset否则默认选 DataFrame 就行别给自己加戏。1.3 任务调度机制一条 SQL 是怎么变成集群任务的这个机制理解了面试和排查问题就都有了底。我给你串一条完整链路你在 SparkSession 里执行一条 SQL首先进入 Catalyst 优化器。它会把 SQL 解析成逻辑计划做谓词下推、列裁剪、常量折叠等优化再转换成物理计划。物理计划会告诉 Spark 该用哪种方式去执行这个查询比如是广播连接还是排序合并连接。接下来是 DAGScheduler 的活。它把整个作业根据宽依赖划分成多个 Stage。窄依赖比如 map、filter是指父 RDD 的每个分区最多被子 RDD 的一个分区使用不需要 shuffle可以尽量合并宽依赖比如 groupByKey、join需要 shuffle 数据必须作为 Stage 的边界。划分好 Stage 之后DAGScheduler 把每个 Stage 拆成一批 Task交给 TaskSchedulerTaskScheduler 再把 Task 分发到集群的 Executor 上执行。实际排查问题的时候只要看到某个 Stage 卡了很久你就要意识到这里大概率发生了 shuffle 或者数据倾斜。这个判断逻辑是通用的往下我会详细讲怎么通过 Spark UI 去精确定位。1.4 Catalyst 优化器容易被忽略的几个优化点Catalyst 是个基于规则的查询优化器它做的很多优化你在日常写 SQL 时已经在享受了只是没有察觉。举个例子谓词下推select * from A join B on A.id B.id where A.dt 2025-01-01Catalyst 会先把 A 表的过滤条件下推到读取数据之后立即执行而不是先 join 完再过滤。这个优化能极大减少参与 join 的数据量尤其是大表和小表 join 的场景。列裁剪也是类似逻辑如果你只 select 三列Spark 读取文件时就不会把整行都加载进来。配合 Parquet 这样的列式存储格式效果非常明显。我自己在调优时会先看 SQL 的执行计划explain命令确认这些优化是否真的发生了。如果发现某个过滤条件没被下推通常是写法有问题比如在 join 条件下做了函数运算导致谓词无法下推。2. 应用场景分析Spark 到底在哪些岗位干活2.1 数仓 ETLSpark SQL 与 Hive 生态的黄金组合数仓 ETL 是目前 Spark 应用最密集的场景。一个典型的离线数仓日任务通常包括几十甚至上百张表的清洗、转换、汇总。用 Spark SQL 跑这些任务比传统的 Hive on MapReduce 在速度上的提升是数量级的特别是复杂多级 join 和聚合体验差异非常明显。在生产环境里最常见的技术组合就是 Spark SQL 加 Hive 元数据服务。Spark 可以无缝读取 Hive 的元数据直接操作 Hive 表甚至可以把 Spark SQL 计算完的结果写回 Hive 表。这样做的好处是你不需要把数仓从 Hive 迁移到另一个平台只需要把计算引擎从 Hive on MapReduce 替换成 Spark SQL业务的改动成本极低。我见过一个实际案例一个团队把原本跑 40 分钟的 Hive ETL 作业迁移到 Spark SQL 上同样的逻辑跑完只需要 6 分钟。他们并没有做什么花哨的优化只是换了个引擎加了合理的资源参数效果就已经很可观了。如果你团队的数仓还是 Hive on MapReduce 在跑强烈建议考虑这个迁移路径。2.2 流式处理Structured Streaming 的适用范围Spark Streaming 的老 API 是基于微批的 DStream现在已经不再推荐。真正值得关注的是 Structured Streaming它把流当成一张无限追加的表来对待你可以用和批处理完全一样的 DataFrame API 去写流处理逻辑。我的经验是Structured Streaming 适合秒级到分钟级延迟要求的场景比如实时大盘、实时风控特征计算、日志实时清洗入湖。它已经能处理事件时间、水印、延迟数据这些复杂语义稳定性也够用。但如果你的场景需要毫秒级延迟比如高频交易那还是建议直接上 Flink。这是架构选型问题不是谁比谁强的问题而是延迟需求不同。这里有个实战细节Structured Streaming 的 checkpoint 机制非常关键。你要确保 checkpoint 目录是可靠的、持久的存储比如 HDFS 或者云上的对象存储。一旦 checkpoint 目录损坏或丢失整个流式作业的精确一次语义就可能被破坏甚至出现重复或者丢数据的问题。2.3 与 AI 管线结合从数据预处理到特征工程这几年 AI 大模型火起来之后Spark 在 AI 领域的角色反而更清晰了它依然是训练之前数据预处理和特征工程的主力工具。不管模型训练用 TensorFlow 还是 PyTorch训练样本的生成、清洗、聚合、采样绝大多数还是靠 Spark 来完成因为数据量实在太大了。我见过不少团队把 Spark 的 DataFrame 和 Pandas 混着用小数据用 Pandas大数据用 Spark中间用 PyArrow 转换。这个思路本身没问题但要注意量级边界。一旦数据量超过单机内存Pandas 就会非常吃力甚至崩溃这时候就必须把流水线迁到 Spark 上。另外现在也出现了一些专门面向 AI 场景的软硬件一体化方案比如搭载 GPU 的 AI 一体机里往往也会把 Spark 作为数据分析和预处理引擎配套部署这说明在大数据架构里Spark 和数据预处理仍然是一对固定搭配。2.4 外部数据源集成Spark 读取 Redis 的正确姿势除了读 HDFS、Hive、对象存储这些常规数据源Spark 读取 Redis 也是一个常被问到的需求。常见的用法是把 Redis 里的维度表数据作为 lookup在 Spark 作业里关联用户画像、商品属性等实时变化的维度信息。方案上主要走 spark-redis 这个连接器它支持把 RDD、DataFrame 写入 Redis也支持按 key 读取 Redis 数据。但我要提醒一点Spark 本质上是一个批量处理引擎你千万不要在 Spark 的算子内部去逐条调用 Redis 的 get 命令。那样会变成海量的小请求打到 Redis性能极差还可能直接把 Redis 打挂我之前就见过这样的生产事故。正确的做法是要么用 spark-redis 连接器做批量读取要么把 Redis 的维度数据先一次性拉成 DataFrame然后做广播变量或者 join。这样可以避免海量的小网络请求性能会好得多。至于说 Redis 的数据量特别大、广播不下的时候就需要做更细粒度的缓存策略了这个是另一个话题后面有机会再展开。3. 性能优化实操让每一台服务器都物有所值3.1 资源参数怎么定先算好 Executor、核心数和内存性能优化里最容易立竿见影的就是资源参数配置。很多新手一上来就问“多少个 Executor 合适”这个问题没法直接回答你得先看集群节点配置再倒推。我拿一个典型场景举例假设你有 5 台数据节点每台机器 32 核、128GB 内存。一般建议每个 Executor 分配 4 到 5 个核这样每个 Executor 的并行度足够又不会因为核心数太多导致线程竞争加剧。如果每个 Executor 5 个核那单机可以跑 6 个 Executor32 除以 5取整留余量剩下 2 个核留给操作系统和节点上的其他进程。内存方面先算单机内存 128GB。要给操作系统预留内存一般留 20% 左右也就是 25GB 左右。剩下 100GB 左右如果单机跑 6 个 Executor每个 Executor 大约 16GB。这里要注意一个参数Executor 的堆内存之外还需要额外开销在 YARN 模式下对应spark.yarn.executor.memoryOverhead默认是 executor 内存的 10%但最小是 384MB。如果 Executor 内存 16GB那 overhead 大约 1.6GB实际分配时要一起考虑。我给出一个比较稳妥的起步配置你可以根据自己的集群调整spark.executor.memory16g spark.executor.cores5 spark.executor.instances30 spark.dynamicAllocation.enabledtrue spark.sql.shuffle.partitions200这个配置的核心逻辑是Executor 数量乘以每个 Executor 的核心数决定了整个作业的并行度上限。spark.sql.shuffle.partitions默认是 200如果你的数据量很大200 个分区可能不够需要调大如果数据量很小200 又多了会造成大量空任务。这个参数需要根据 shuffle 数据量反复试没有一个固定的万能值。3.2 序列化方式与存储格式看似小细节影响大性能序列化方式直接影响 Spark 作业的 CPU 和内存开销。Spark 默认使用 Java 序列化但它的问题很明显序列化后的数据体积大序列化过程本身也慢。换成 Kryo 序列化之后数据体积和序列化时间通常能减少一半以上尤其是复杂对象场景提升更明显。启用 Kryo 的方式很直接在提交作业前设置spark.serializerorg.apache.spark.serializer.KryoSerializer然后把你注册的自定义类写在--conf里或者在代码里调用conf.registerKryoClasses。如果不注册Kryo 也能工作但会给每个类写完整的类名序列化出来的数据就更大所以还是要养成注册类的好习惯。存储格式方面我强烈建议在能选择的情况下优先用 Parquet。它是列式存储配合 Catalyst 的列裁剪读数据时可以只读需要的列极大地减少了 I/O。而且 Parquet 自带压缩压缩比高读取性能也好。我见过有团队还在用 CSV 存中间结果同一份作业换成 Parquet 后查询时间直接缩短了一半以上这个优化做起来几乎零成本。3.3 分区策略与数据倾斜先把数据分均匀再说并行数据倾斜是大数据场景里最常见的性能杀手几乎每个做 Spark 开发的人都会遇到。它的症状很典型某个 Stage 里大部分 task 很快就跑完了但有一两个 task 卡很久拖得整个作业迟迟无法结束。根本原因是某些 key 的数据量远远大于其他 key比如电商场景里的热门店铺、社交场景里的头部用户。这些 key 所在的分区数据量巨大单个 task 处理不过来别的 task 又在空转。我处理数据倾斜通常按以下顺序排查和解决先确定是哪个 Stage 出了问题从 Spark UI 的 Stage 详情里看每个 task 的处理时间和 shuffle 数据量。确认倾斜 key 是什么写一个简单的 group by 查询按照数据量排序你可以根据实际场景去确定热点 key。根据场景选择合适的方案比如过滤掉异常数据、增加分区数、加盐打散、广播小表等。其中加盐Salting是最常用的硬核方案思路就是给热点 key 加上随机后缀让它们被分散到多个分区去处理完成 join 或聚合后再去掉后缀。后面我会用一个完整的案例演示怎么做。3.4 Join 优化广播连接和排序合并连接怎么选Join 是 Spark 作业里最消耗资源的操作之一选择正确的 join 策略对性能的影响极大。广播连接Broadcast Join适合大表和小表 join 的场景。Spark 会把小表的数据广播到每个 Executor 上避免 shuffle性能极快。默认的广播阈值是 10MB由spark.sql.autoBroadcastJoinThreshold控制。如果小表只有几十 MB你可以在 SQL 的注释里强制指定广播比如SELECT /* BROADCAST(dim) */ fact.*, dim.name FROM fact JOIN dim ON fact.dim_key dim.id但要注意广播变量会复制到每个 Executor 的内存里。如果小表已经很大比如几百 MB广播的开销会非常大甚至可能导致 Executor 内存溢出。所以广播连接不是越大越好需要根据 Executor 的内存和集群规模综合判断。排序合并连接Sort Merge Join是默认的大表 join 策略。它的核心是先对两个表按照 join key 排序然后做归并连接这个过程需要 shuffle。优化这个场景的关键在于一是避免在 join 字段上做函数运算二是在满足业务的前提下尽量先过滤再 join缩小参与计算的数据集。另外一个常见的坑是两个大表 join 后只取少数行正确的做法是先聚合出结果再 join而不是 join 完再筛选。4. 内存管理从 JVM 结构到 OOM 排查实录4.1 Spark 统一内存模型一张图记住内存分配的底层逻辑Spark 的内存管理在 2.0 之后采用了统一内存模型核心是把 Executor 的堆内存分成三块Reserved Memory保留内存固定 300MB、User Memory用户内存用于存储用户自定义的数据结构和 RDD 转换产生的对象和 Spark MemorySpark 内部统一管理执行和存储两部分的动态内存。其中 Spark Memory 默认占堆内存的 60%由spark.memory.fraction控制。这 60% 里执行内存shuffle、join、聚合时使用和存储内存缓存 RDD、广播变量等可以互相借用由spark.memory.storageFraction决定初始划分比例默认 50%。理解这个模型对排查 OOM 很有用。比如当你执行大批量 shuffle 时执行内存可能占用大量空间此时如果广播变量也很大存储内存不足两者就会相互争抢最后可能引发频繁 GC 或者直接 OOM。遇到这种情况你可以提高 Executor 内存、降低缓存的数据量或者把部分内存挂到堆外。4.2 四种最常见的 OOM 场景与对策第一个是 Driver 端 OOM。常见原因是collect()把大量数据拉回 Driver或者数据量太大的 DataFrame 被转成了 Pandas。解决办法很简单尽量不要collect()大数据集要么转成聚合结果要么分页拉取要么通过foreachPartition写入外部系统。第二个是 Executor 端 OOM。多半发生在 shuffle 过程中某个分区数据量过大或者合并大量小文件时产生大量对象。应对思路是调整spark.sql.shuffle.partitions让分区更均匀或者适当提高 Executor 内存更重要的是从源头避免数据倾斜。第三个是堆外内存 OOM。这通常发生在使用了堆外缓存、或者网络缓冲区压力过大的时候。解决办法是提高spark.yarn.executor.memoryOverhead或spark.memory.offHeap.size。这类 OOM 在日志里往往不是明显的OutOfMemoryError而是作业失败、Executor 被容器杀掉排查时要多留个心眼。第四个是广播变量 OOM。广播数据量超过 Executor 内存能承受的上限时会导致 Executor 内存迅速膨胀。如果你发现广播变量很大就别硬用广播连接改成普通 join或者把广播数据做成外部存储的查找表。4.3 从 Spark UI 一眼锁定瓶颈Spark UI 是我排查任何性能问题的第一站也是最有用的一站。打开一个作业的 UI重点看三个页面Executors 页面看每个 Executor 的内存使用情况、GC 时间和 Shuffle 读写量。如果某个 Executor 的 GC 时间明显高于其他节点说明它处理的数据量有倾斜或者堆内存压力过大。SQL 页面它展示了完整执行计划并且精确标出了每个操作消耗的时间。你很快就能看到是哪个 join、哪个聚合、哪个读取操作占据了绝大部分时间。我曾经仅凭这个页面就定位到一个用正则表达式处理大字段的 JSON 解析它占了整个作业 70% 的耗时优化思路瞬间就有了。Stage 页面看 Task 的时间曲线。如果出现某几个 task 时间很长而其他很短就是数据倾斜如果所有 task 的调度延迟很高说明集群资源不足有作业在抢占如果某个 Stage 的 Shuffle Read 数据量特别大就要考虑是不是 join 或 group by 导致的数据爆炸。5. 集群搭建与部署从零到能跑的经验总结5.1 部署模式怎么选Standalone、YARN 还是 KubernetesSpark 支持的部署模式很多我按实际经验给你列一下适用场景Standalone 模式是 Spark 自带的资源调度搭建最简单适合测试环境和小规模集群也适合刚学习 Spark 的同学用来理解主从架构。我在本地学习时就常用 Standalone 模式配合一个三节点的集群就能模拟完整的分布式计算。YARN 模式适合有 Hadoop 生态的团队。因为 YARN 是统一的资源调度器可以同时跑 Spark、Flink、MapReduce 等多种计算框架资源利用率高。你的集群如果已经运行了 Hive那基本上就是 YARN 模式。Kubernetes 模式是现在的新趋势。优点是资源隔离好、弹性伸缩快、环境一致性高适合容器化基础设施完善的团队。缺点是运维复杂度高尤其是状态管理和网络配置需要团队有较强的容器平台能力。5.2 安装与配置的详细步骤和避坑点我以最常见的 Spark on YARN 为例说一下安装的关键步骤下载 Spark 安装包时要特别注意版本兼容性。Spark 3.x 和某些 Hadoop 版本之间可能存在协议不兼容建议先看官方文档确认你正在使用的 Hadoop 版本再选择对应的 Spark 版本。Java 版本也要匹配Spark 3.x 通常支持 Java 8/11/17但不同小版本可能有差异稳妥起见选定主版本后保持环境一致。配置环境变量时重点设置SPARK_HOME和JAVA_HOME把$SPARK_HOME/bin加进 PATH。如果集群是部署模式需要在每台节点都做同样的配置。修改 Spark 的核心配置文件spark-defaults.conf时把公共参数比如 Driver/Executor 内存、序列化方式、shuffle 分区数写入这个文件这样作业提交时就有了合理的默认值。一个非常常见的坑是集群里同时存在多个 Hadoop 版本Spark 运行时找不到正确的 Native 库或 classpath。这个问题通常表现为启动时报NoClassDefFoundError。解决的思路是确认HADOOP_CONF_DIR和HADOOP_HOME指向正确的路径必要时把 HDFS 配置文件软链到 Spark 的 conf 目录下。5.3 连接云上存储的几个关键参数现在很多 Spark 作业不是读 HDFS而是直接读对象存储比如 S3、OSS 或者 GCS。连接这些存储时有几个参数很容易被忽略一是访问密钥的配置建议放在配置文件里授权不要写在代码中二是连接超时和重试参数对象存储在出现异常时合理的重试次数能在很大程度上解决网络抖动的问题三是小文件问题对象存储对小文件的读写性能很差你会发现大量文件在读取时打开了非常多的小分片导致效率非常低。解决办法是在写入后做一次文件合并或者控制好写入时的分区数量避免产生过多碎片文件。6. 常见问题速查表生产线上的踩坑清单我整理了一张问题排查速查表这些都是我在生产环境里真实遇到过、也经常在团队里被问到的场景。你可以把它收藏起来遇到类似问题时快速对照。现象可能原因检查方向常见解法作业运行到某个 Stage 长时间不结束数据倾斜Spark UI 中 Stage 详情查看各 Task 耗时加盐、重分区、过滤极端 keyExecutor 反复被杀报 container 内存超限堆外内存或对象满溢Executor 日志YARN 容器监控提高 memoryOverhead减少缓存优化代码大量小文件导致读写慢分区数过多输出时未合并看输出文件数和大小coalesce 或 repartition 控制设置动态分区合并Shuffle 读数据量巨大网络成为瓶颈join/group by 数据膨胀SQL 执行计划Shuffle Read 指标过滤下推、联 key 做加盐、更换 join 策略GC 频繁、任务曲线锯齿状Executor 内存偏小对象过多Executor 页面 GC 时间提高 Executor 内存Kryo 序列化减少 shuffle 数据任务调度延迟高等待时间长集群资源不足有作业抢占ResourceManager 页面队列使用率调整队列配额错峰运行动态资源分配Spark 写 Redis 极慢或打挂 Redis在算子内逐条调用 Redis检查代码是否有循环请求改用 spark-redis 批量读写或批量写入作业使用 Hive 元数据报表不存在Spark 用户没有 Hive 库的访问权限或元数据版本不匹配检查用户权限、Hive Metastore 版本配置正确的 HMS 连接并检查用户认证这张表不是万能的但能覆盖大部分 Spark 生产问题的排查入口。我个人的体会是大多性能问题走到最后都指向两个根源要么是数据分布不均匀要么是资源配置和数据结构不匹配。把这两条主线记在心里排查效率会高很多。7. 实战案例一次生产作业从 22 分钟到 4 分钟的调优过程前面讲了那么多理论和方法最后我分享一个完整的实战案例。这个作业是典型的双表 join 聚合一张订单事实表一张用户维度表按照用户维度统计每个用户最近的订单量。原先作业跑 22 分钟优化后只需要 4 分钟。第一步是定位问题。打开 Spark UI发现 Stage 3 有大量 task 几秒就完成了但是有 3 个 task 跑了接近 20 分钟。明显是数据倾斜。我再做一些进一步分析发现用户表中有一个头部用户比如一个测试账号或者超级大客户它的订单量占了全量订单的 30% 左右这就是倾斜的根源。第二步是确定方案。因为用户维度表整体不大大约 10GB不适合直接广播。我先做 join 前的加盐把订单表中该热点的用户 ID 加上随机后缀把用户维度表的数据拆成多份然后再 join。整个过程需要做两次聚合第一次按加盐后的 key join第二次去掉盐值按真实用户 ID 做最终的聚合。代码用 PySpark 的 DataFrame API 写关键逻辑如下from pyspark.sql import functions as F # 订单表中给倾斜key加随机后缀 salted_orders orders.withColumn( join_key, F.when( F.col(user_id) hot_user, F.concat(F.col(user_id), F.lit(_), (F.rand() * 10).cast(int)) ).otherwise(F.col(user_id)) ) # 用户维度表对倾斜key行做多份复制后缀范围 0-9 salted_users users.withColumn( join_key, F.when( F.col(user_id) hot_user, F.explode(F.array([F.lit(fhot_user_{i}) for i in range(10)])) ).otherwise(F.col(user_id)) ) # 第一轮 join在加盐后的key上进行 joined salted_orders.join(salted_users, join_key, left) # 第二轮聚合先按盐值key聚合再按真实user_id聚合 result ( joined.groupBy(join_key, user_id) .count() .groupBy(user_id) .agg(F.sum(count).alias(order_cnt)) )第三步是验证效果。调整后加盐的那批热数据被分散到 10 个任务里处理不再压在某一个 task 上整个作业的运行时间从 22 分钟降到了 4 分钟。集群负载也均衡了不再出现某台机器 CPU 打满、其他机器空转的情况。这个案例里最关键的一步不是写加盐代码而是通过 Spark UI 准确定位数据倾斜的阶段和 key然后再有针对性地设计盐值范围和分布。盐值范围要根据热点数据的量级来定如果热点数据很大10 份不够就取 50 份以最终任务均衡为准。加盐虽然会导致一次多余的聚合但相比几十倍的性能差距这个代价完全可以接受。另一个补充案例是关于 Spark 读取 Redis 的。之前有个实时特征作业为了取用户画像数据在 map 算子里面逐条调用 Redisget结果 Redis 的 CPU 被打到 90%作业也经常超时。后来我改成每天把画像表同步成 Hive 表用 Spark 直接 join这个场景就消失了。如果你的场景确实需要实时读取 Redis至少要用spark-redis的批量接口并控制好并发和 key 的分布。8. 写在最后Spark 这个生态覆盖的面太广了没有人能说自己全都精通。我这几年的经验总结下来其实就是三条第一理解任务调度和内存模型这是排错的基础第二任何性能优化都要先基于 Spark UI 和数据特点做诊断不要凭空猜第三把常见问题沉淀成排查清单团队里每个人都用得上。这套系列后面我再接着写一些深度内容比如 Spark 与数据湖的集成、Structured Streaming 的 exactly-once 实践、以及大集群的资源调优策略。如果你在实际操作中遇到好玩的坑也欢迎随时交流。毕竟大数据这条路一个人走容易踩坑互相分享才能走得更远。
返回列表