ARTICLE DETAIL

资讯详情

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

SparkStreaming+Kafka参数调优实战:从Direct模式到Offset管理,避坑指南

SparkStreaming+Kafka参数调优实战:从Direct模式到Offset管理,避坑指南 做流式计算的人几乎都绕不开 SparkStreaming 和 Kafka 这对组合。搞懂它们之间的参数配置绝对比跑通一个 demo 要难得多——很多任务刚上线时看着正常跑一周后就开始出各种幺蛾子消费延迟越积越高、重启后重复消费、时不时再来个 OOM。这篇文章把我这些年调 SparkStreaming Kafka 参数的经验一次性整理出来包含每个参数背后的原理、推荐值、线上踩过的坑以及一份能直接抄的配置模版。适合正在用 SparkStreaming 做实时计算、或者准备从批处理转流处理的同学参考。1. 先搞懂 Kafka 和 SparkStreaming 是怎么连起来的1.1 两种集成方式Receiver 与 Direct先说结论现在基本只用 Direct 模式。但要理解参数配置必须先知道为什么会有两种模式。早期 Spark Streaming 用 Receiver 方式接入 Kafka。Receiver 作为一个长期运行的 Task 常驻在 Executor 上不断从 Kafka 拉取数据把数据放到 Executor 的 Block Manager 里。Spark Streaming 再按照固定的批次间隔把 Block 中的数据封装成 RDD 去计算。这个模式的 offset 是自动保存到 ZooKeeper 的你不需要手动管理。听起来很省事但问题也出在这数据先落在 Executor 内存里如果一批数据还没来得及被处理Executor 就崩溃了数据就丢了。虽然可以开启 WAL 写 HDFS 来减少丢失但 WAL 的写入开销非常大会严重拉低吞吐。此外Receiver 模式和 Spark 的批次调度是割裂的上游拉取速率和下游消费速率很难匹配一旦数据峰值到来要么积压要么直接把 Executor 内存打爆。Direct 模式是 Kafka 0.10 之后主推的集成方式。DirectStream 不再需要 Receiver而是由 Driver 在每次 batch 开始时通过 Kafka Consumer API 拿到每个分区的最新 offset再统一计算本次要读取的 offset 范围Executor 直接从 Kafka broker 拉数据不需要经过 Block Manager。数据不落 Executor 内存缓冲offset 也可以由任务自己管理真正做到了“一个批次对应一个准确的数据范围”。1.2 为什么现在都在用 Direct 模式Direct 模式最大的优点有两个。第一并行度模型清晰RDD 的分区数和 Kafka 的分区数一一对应你可以直接通过调整 Kafka 分区数来控制 SparkStreaming 的并行度。第二offset 管理可控由于 Driver 知道每个分区实际消费到哪个位置可以将 offset 手动提交到 Kafka 或外部存储配合幂等输出能实现接近精确一次exactly-once的效果。从参数配置角度来说Receiver 模式里要考虑 spark.streaming.receiver.writeAheadLog.enable 这类参数而 Direct 模式更关注的是 Kafka Consumer 参数、Spark 批次参数和 offset 管理参数。这篇文章后面讲的内容全部以 Direct 模式为前提。如果你还在维护老项目里的 Receiver 代码我建议尽早迁移不只是因为参数更好调代码的可维护性也完全不一样。2. 必须吃透的核心参数先从 Kafka 侧说起SparkStreaming 的 DirectStream 底层调用的就是 Kafka 原生 Consumer API所以 Kafka Consumer 的参数会直接作用在 SparkStreaming 任务上。很多人排错时只盯着 Spark 日志经常忽略 Kafka 这边的参数导致同一个问题反复出现。这一节把 Kafka 侧常用参数按功能分组过一遍。2.1 消费组与 offset 参数group.id 是消费组 ID多个 consumer 实例要协同消费同一个 topic 时必须使用相同的 group.id。SparkStreaming 任务重启时如果 group.id 不变Kafka 会根据已提交的 offset 继续消费如果 group.id 变了就相当于一个全新的消费组会在 auto.offset.reset 指定的位置开始消费。所以 group.id 一定要固定并且命名尽量带上业务含义比如 dwd_user_order_rt方便监控和排查。enable.auto.commit 和 auto.commit.interval.ms 控制“是否自动提交 offset”。在 Direct 模式下我强烈建议把 enable.auto.commit 设为 false因为自动提交的时机和 Spark 批次处理完成的时机完全对不上。假设你的 batch 是 10 秒但 Kafka 每 5 秒自动提交一次那完全可能出现一种情况某一批数据还没处理完offset 已经提交了这时任务崩溃重启后就会从已提交的位置开始消费中间那些正在处理的消息就丢了。反过来如果数据处理完了但 offset 还没提交重启后会重复消费一部分数据。自动提交带来的问题就是提交时机不可控所以大多数生产方案都会关掉它改成手动提交而且要等整个 batch 处理完成、输出写成功之后再提交。auto.offset.reset 决定在 Kafka 中没有已提交 offset 或者提交的 offset 已经失效时consumer 从哪个位置开始消费。取值有 earliest、latest、none。earliest 从最早消息开始消费latest 只消费新消息如果设置 none 且没有提交过 offset直接抛异常。对于实时数仓场景我一般建议 earliest这样 topic 里有历史数据时可以通过重置 offset 回溯。但如果你是新上线任务topic 里旧数据已经不重要就用 latest否则第一次启动会把积压几个月的数据全读一遍直接把下游打挂。这里没有绝对正确必须结合业务需求选。2.2 拉取与网络参数fetch.min.bytes 表示一次拉取请求最少需要返回的数据量。如果 broker 端累积的数据不足这个值fetch 请求会等待 fetch.max.wait.ms 再返回。调大 fetch.min.bytes 可以减少请求次数、提升吞吐但会增加延迟。在 SparkStreaming 场景下批次本身就是批量拉取所以默认的 1 字节一般就够用不需要刻意调大。fetch.max.wait.ms 默认 500 毫秒如果你的 topic 数据密度低可以适当调大减少空轮询但不建议超过 1000否则批次内延迟会变大。max.partition.fetch.bytes 控制每个分区一次最多拉取多少字节默认 1MB。如果你的单条消息很大比如超过 1MB或者希望减少拉取次数可以调大到 5MB 甚至 10MB。但这个值千万别瞎调因为 DirectStream 会把拉取的数据作为 RDD 分区缓存到 Executor 内存单次拉取量越大RDD 分区的内存占用就越高调太猛很容易 OOM。尤其是 1MB 是个很常见的坑当消息体超过 1MB 时如果 max.partition.fetch.bytes 不调整consumer 会一直拿不到数据表现为“任务正常但消费不到新消息”。receive.buffer.bytes 和 send.buffer.bytes 是 socket 层缓冲区默认 64KB。在万兆网卡或跨机房传输场景可以调到 256KB 或 512KB减少 TCP 读写次数。但这两个参数只影响网络传输层如果瓶颈不在 Kafka 到 Spark 的带宽上调了也没用。2.3 会话与心跳参数session.timeout.ms 和 heartbeat.interval.ms 是配套使用的。session.timeout.ms 表示 consumer 在多少秒内没有发送心跳就被判定为挂掉触发 rebalance。heartbeat.interval.ms 是心跳发送间隔通常设置成 session.timeout.ms 的三分之一。在 SparkStreaming 任务里一个 Executor 线程可能在长时间处理某个 batch导致暂时没有空闲发心跳。如果把 session.timeout.ms 设得太小就会出现“任务明明在跑却被踢出消费组”的诡异问题。所以建议把 session.timeout.ms 调到 25 秒到 30 秒heartbeat.interval.ms 相应调到 5 秒左右。max.poll.interval.ms 比 session.timeout.ms 更容易踩坑。它表示 consumer 两次 poll 之间的最大间隔超过这个时间没有 pollconsumer 就被认为消费能力不足会主动离开消费组并触发 rebalance。在 SparkStreaming 里一个 batch 处理期间consumer 线程是不会去 poll 的如果 batch 处理时间超过 max.poll.interval.ms就会触发 rebalance。默认值是 5 分钟但实时任务一个 batch 跑 5 分钟以上并不罕见尤其是大数据量 join 或者写 HBase 抖动的时候。这个参数一定要根据你任务的最长 batch 耗时来设置建议设置成 max(batchDuration * 3, 300000) 毫秒保守可以到 10 分钟甚至更长。max.poll.records 控制一次 poll 最多返回多少条消息默认 500。如果单条消息处理成本高可以把这个值调小减轻单批次的处理压力如果消息都是小体积且逻辑简单可以调大减少 poll 次数。3. SparkStreaming 侧参数从批次到并行度Kafka 侧参数解决的是“能不能正确消费”SparkStreaming 侧参数解决的是“消费完了能不能及时处理完”。很多任务延迟高不是 Kafka 拉不过来而是 Spark 这边批次大小不合理、并行度不够或者内存配置有问题。3.1 batchDuration 怎么定batchDuration 是 SparkStreaming 最核心的参数它决定每隔多长时间生成一个 RDD。这个参数在代码里通过 new StreamingContext(conf, Seconds(10)) 传入。它既不能太小也不能太大。太小的话Spark 的调度开销会占大头如果每个 batch 只有几百条数据大部分时间都浪费在任务调度上太大的话延迟变高数据从产生到可被消费的时间变长实时性变差。我的经验是先看业务延迟要求。如果业务要求秒级batchDuration 可以设 2 秒或 5 秒但消费的调度资源更高如果接受 10 秒到 30 秒延迟设 10 秒或 15 秒通常更稳。其次看单个 batch 的处理耗时理想情况是处理耗时小于 batchDuration 的 70%。如果你发现平均处理耗时已经接近 batchDuration说明系统过载了需要调大 batchDuration、增加资源或者优化处理逻辑。千万不要抱着“batchDuration 越小越实时”的想法对一个过载的系统调小 batchDuration 只会雪上加霜。还有一个隐蔽问题batchDuration 太小会导致很多 batch 只拉到很少的数据尤其是数据稀疏的 topic。这些空 batch 不仅浪费调度资源下游写 HDFS 时还会产生大量小文件后续查询性能会很难受。3.2 spark.streaming.* 参数详解spark.streaming.backpressure.enabled 和 spark.streaming.backpressure.initialRate 是背压参数。开启后系统会根据当前 batch 处理速度动态调整最大消费速率防止数据积压导致系统崩溃。对于流量波动很大的业务强烈建议开启。同时要设置 spark.streaming.kafka.maxRatePerPartition它表示每个 Kafka 分区每秒最多消费多少条消息。背压开启后实际消费速率会在 initialRate 和 maxRatePerPartition 之间动态调整。我一般把 maxRatePerPartition 设置为略高于正常峰值流量的值比如峰值的 1.2 到 1.5 倍这样既能保护系统又不会在流量高峰时把数据卡在 Kafka 里。spark.streaming.kafka.maxRetries 在 Direct 模式下如果某个分区拉取失败最多重试多少次默认 3。如果 Kafka 集群本身不太稳定可以适当调大但更建议从 Kafka 集群和网络层面解决。spark.streaming.stopGracefullyOnShutdown 设置为 true才可以在停止任务时等待正在处理的 batch 完成。关闭的话任务被 kill 时正在处理的数据会直接中断很可能丢数据。虽然 offset 手动提交会兜底但还是建议打开配合监控做优雅下线。3.3 并行度与资源参数Direct 模式下每个 Kafka 分区对应 RDD 的一个分区所以并行度首先受 Kafka 分区数限制。如果你想提升消费并行度首选方案是增加 Kafka 分区数其次是在算子内部通过 repartition 重分区。但 repartition 会引入 shuffle有额外开销不需要重分区的场景不要随便加。资源参数方面executor 数量和每个 executor 的核数决定了整体并行能力。我一般建议每个 executor 2 到 5 个核避免单 executor 内线程太多导致 CPU 争抢。Executor 内存要考虑每个 batch 会缓存输入数据一般设置 spark.executor.memory 和 spark.executor.memoryOverhead 时要预留出至少两倍输入数据量的空间。如果数据量大建议打开 spark.memory.offHeap.enabled并设置 spark.memory.offHeap.size把一部分缓存放到堆外减少 GC 压力。但注意堆外内存也不是无上限的操作系统本身还要留余量调太大整台机器都可能被拖垮。4. 一套能直接抄的配置模版光讲参数不给出能跑的例子等于白讲。下面给出一份基于 Direct 模式的 SparkStreaming Kafka 参数配置参考分为代码层和提交脚本层。4.1 代码层面的参数设置示例这里用 Scala 写一个简单的 DirectStream 示例把 Kafka Consumer 参数放进 Map再传给 createDirectStream。其中最重要的是 enable.auto.commitfalse以及手动提交 offset 的逻辑。val kafkaParams Map[String, Object]( bootstrap.servers - kafka1:9092,kafka2:9092,kafka3:9092, key.deserializer - org.apache.kafka.common.serialization.StringDeserializer, value.deserializer - org.apache.kafka.common.serialization.StringDeserializer, group.id - dwd_user_order_rt, auto.offset.reset - latest, enable.auto.commit - (false: java.lang.Boolean), session.timeout.ms - 30000, heartbeat.interval.ms - 5000, max.poll.interval.ms - 600000, max.poll.records - 1000, fetch.max.wait.ms - 500 ) val topics Array(dwd_user_order) val stream KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) )注意几个细节。key.deserializer 和 value.deserializer 必须设置否则会直接报 ClassNotFound。enable.auto.commit 必须写成 java.lang.Boolean 类型因为 Map 的类型是 Map[String, Object]直接写 false 会当成 Scala Boolean运行时类型不匹配这也是我踩过的坑。session.timeout.ms 和 max.poll.interval.ms 设得宽松是为了给批次处理留足时间避免 rebalance 风暴。4.2 提交脚本中的 Spark 参数代码层配置好后spark-submit 里的资源配置同样关键。下面是一份非常通用的提交参数spark-submit \ --class com.example.StreamingJob \ --master yarn \ --deploy-mode cluster \ --driver-memory 4g \ --executor-memory 8g \ --executor-cores 4 \ --num-executors 20 \ --conf spark.streaming.backpressure.enabledtrue \ --conf spark.streaming.backpressure.initialRate5000 \ --conf spark.streaming.kafka.maxRatePerPartition10000 \ --conf spark.streaming.stopGracefullyOnShutdowntrue \ --conf spark.serializerorg.apache.spark.serializer.KryoSerializer \ --conf spark.memory.offHeap.enabledtrue \ --conf spark.memory.offHeap.size2g \ --conf spark.task.maxFailures8 \ my-streaming.jarexecutor 数量和 cores 要和 Kafka 分区数匹配。理想情况下executor 总核数大于等于 Kafka 分区数这样每个分区都能被并行处理。比如 topic 有 30 个分区executor-cores 是 4那至少需要 8 个 executor 才能让所有分区并行。如果总核数小于分区数有些分区会被排队处理延迟升高。当然总核数也不能比分区数多太多多出来的核心会空转还增加调度开销。这里提醒一句KryoSerializer 和实时任务很搭默认 Java 序列化性能差、对象大、GC 压力高。但是 Kryo 对某些 Scala 集合类的兼容性不好如果遇到序列化异常需要手动注册类或者调整配置。4.3 offset 管理的正确姿势关闭自动提交后offset 到底存在哪最简单的方案是存 Kafka 内部主题 __consumer_offsets通过 commitAsync 提交。优点是实现简单缺点是不够灵活下游要精确统计消费到哪或者在多个 group 之间共享 offset 时就不方便了。更稳妥的做法是把 offset 存到外部存储比如 HBase 或 MySQL并且在同一个事务里写业务结果和 offset这样能做到真正意义上的“至少一次”或“精确一次”。具体思路是处理逻辑中先把业务结果写到一张带唯一键的表同时把该批次的 offsetRanges 也写入同一个事务。下次任务启动时从这张表恢复 offset并通过唯一键去重下游就不会因为重复消费产生脏数据。生产环境里因为 offset 丢失导致重复消费的案例太多了大部分是没把 offset 管理当回事。5. 线上问题排查实录与调优心得参数配置最终都要落到线上问题的解决上。挑几个最常见的场景按“现象—原因—解法”来说。5.1 消息消费延迟高这是被问得最多的问题。现象是 Kafka 消费组 Lag 持续上涨SparkStreaming 任务每个 batch 的处理时间接近甚至超过 batchDuration。首先要确认是“拉取不过来”还是“处理不过来”。看 Spark UI 每个 batch 的 Scheduling Delay如果调度延迟特别大说明资源不够如果 task 计算耗时大说明处理逻辑是瓶颈。如果 Kafka 端本身有大量积压而 Spark 资源已经加不上去就先开启背压降低消费速率防止雪崩。然后评估是扩容 topic 分区数还是增加 executor。分区数扩容对 Kafka 侧提升明显但扩容后需要重启 SparkStreaming 任务重新分配分区而且对消息顺序性要求高的场景要谨慎。另一个容易忽略的点是输出端写入性能比如写 HBase 时 batch 大小是否合理、连接池是否满了。很多时候 Spark 处理并不慢慢在下游存储。5.2 重复消费与丢数据重复消费绝大多数和 offset 提交时机有关。自动提交一定会有重复消费的可能解决办法就是关闭自动提交改成业务处理成功后再提交同时在下游做幂等。丢数据要分两种一种是真的丢了比如 enable.auto.committrue 时Spark 拉到了数据但 batch 还没处理任务被 killoffset 已经自动提交了这批数据就丢了另一种是看起来丢了实际是处理失败被静默忽略了如果你在 foreachRDD 里写完数据后没有检查返回值也没抛异常数据丢了根本不知道。排查这类问题第一件事是把 Kafka 消费组的 offset 和 Spark 实际处理的数据量做对比。我习惯在每次 batch 提交 offset 时打印 offsetRanges 的 fromOffset 和 untilOffset并记录一条“batch 处理了多少条、输出成功多少条”的日志。数据对不上时第一时间就能判断是拉取少了还是输出少了。5.3 频繁 rebalance 与 OOM频繁 rebalance 通常是 session.timeout.ms、heartbeat.interval.ms、max.poll.interval.ms 设置不合理导致的。尤其是有些 batch 处理时间超过 max.poll.interval.msconsumer 会主动离开组。这个问题在 SparkStreaming 里很典型因为 Spark 的 batch 处理时间经常被写存储的抖动放大。建议把 max.poll.interval.ms 调到 10 分钟以上并配合监控观察 rebalance 次数。另外同一个 group 下有多个 SparkStreaming 应用实例也会导致 rebalance 混乱一定要确认线上没有重复提交同一个应用。OOM 要分 driver 和 executor。Driver OOM 常见于 foreachRDD 里把大量数据 collect 回 driver或者 checkpoint 写得太频繁。Executor OOM 常见于输入数据量过大、缓存内存不足。解决办法除了调大资源还要减少单个批次的数据量比如调大 batchDuration、调小 max.poll.records、给输入 RDD 做 repartition 分散压力。堆外内存只是一个缓冲不能无脑调大如果堆内堆外一起爆还是要回到数据量和并行度这两个根本问题上。5.4 参数速查表为了方便对照把常用参数的推荐值和使用场景整理成一张表。注意这只是起点值实际要结合数据量、资源、业务要求来调。参数推荐值场景说明batchDuration5s~15s延迟敏感的用 5s吞吐优先用 15senable.auto.commitfalse生产环境必须关auto.offset.resetearliest/latest可回溯用 earliest新任务追增量用 latestsession.timeout.ms25s~30s防止 batch 处理时心跳超时heartbeat.interval.ms5s约为 session.timeout 的 1/6max.poll.interval.ms10min 以上必须大于最长 batch 耗时max.poll.records500~2000单条大消息调小maxRatePerPartition峰值流量的1.2~1.5倍配合背压使用executor 总核数大于等于 Kafka 分区数保证分区能并行消费spark.streaming.stopGracefullyOnShutdowntrue优雅停机不丢批6. 最后再分享一点个人体会参数配置这件事最忌讳的就是从网上复制一份“最优配置”直接套。我踩过很多次坑之后发现每个参数的合理值都和你自己的数据量、资源、业务语义强相关。真正高效的做法是先把 Kafka 侧和 Spark 侧每个参数的含义吃透然后用一套保守的初始配置跑起来通过监控Kafka 的 consumer Lag、Spark UI 的 batch 处理耗时、GC 日志观察瓶颈再针对性地调整一个参数、观察一次效果。一次只动一个变量才能知道是哪个参数起了作用。如果你刚接触这套体系不妨从今天这份模版开始慢慢调出自己的配置组合。最后记住一个原则宁可消费慢一点、batch 大一点也不要让任务频繁失败重启稳定永远是实时任务的第一位。
返回列表