ARTICLE DETAIL

资讯详情

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

Spark Streaming实战:微批模型、DStream与调优避坑指南

Spark Streaming实战:微批模型、DStream与调优避坑指南 做实时数据处理的人迟早会碰到 Spark Streaming 这个名字。它是 Spark 生态里最早承担实时流计算任务的 API和批处理用的 RDD API 共用一套引擎只是把“一次性处理全部数据”变成了“按小批次持续处理不断到达的数据”。简单说Spark Streaming 解决了两个问题一是让 Spark 能处理源源不断的数据流二是提供了 DStream 这样一套专门面向流的编程接口开发人员不用另学一套分布式框架就能把实时业务跑起来。这套 API 成熟稳定在电商、金融、IOT 等场景里沉淀了大量生产案例哪怕现在新一代的 Structured Streaming 已经越来越主流存量系统里大量 Spark Streaming 作业仍然在线上运行很多流式处理的核心思路也都是从它这里来的。这篇文章不会讲太多虚的就按我自己的实操经验把 Spark Streaming 的原理、开发流程、调优手段和踩坑记录完整过一遍。1. 微批模型、DStream 与窗口设计先把原理吃透1.1 微批不是缺陷它是吞吐与容错的取舍很多人一提到 Spark Streaming第一反应就是“它是个假流处理延迟太高”。这个说法一半对一半不对。Spark Streaming 确实不是真正意义上的逐条处理它采用的模型叫微批Micro-Batch数据到达后不是立刻计算而是先攒到一个时间窗口里比如每 5 秒攒一批然后再把这批数据交给 Spark 的批处理引擎去算。这样做最直接的好处是复用了 Spark 最成熟的批处理内核容错、调度、内存管理全都继承自离线计算一套代码既跑批又跑流。代价是延迟的下限被批次间隔卡住了不可能做到毫秒级响应。从我维护生产任务的经验来看微批模型对绝大多数业务场景完全够用。实时大屏、风控规则扫描、优惠券核销、日志清洗这类需求秒级延迟已经能覆盖 95% 以上的情况。真正需要毫秒级响应的场景本质上应该去选专门的流处理引擎而不是在 Spark 上纠结。换一个角度想微批模型反而给了你一个天然的“背压缓冲”如果下游临时抖动数据还能在内存里等一等不会像纯流引擎那样要么阻塞、要么直接丢弃。理解微批模型还有一个很重要的点要记住每批次处理时间是整个 Streaming 作业健康度的核心指标。如果你在 Spark UI 里发现一个批次的处理时间已经逼近甚至超过了批次间隔说明系统进入过载状态接下来会出现任务积压。积压到一定程度Spark 会开始丢弃旧批次这才是真正的数据丢失。1.2 DStream看不见的 RDD 流水线Spark Streaming 的编程入口是 DStream全称 Discretized Stream中文叫离散化数据流。它内部其实是一串不间断的 RDD 序列每个 RDD 对应一个时间片的数据。你可以把 DStream 想成一本电子书每页就是一个 RDD时间是自动翻页器每隔一个批次间隔翻一页。用代码来理解会更直白。假设你要从 Kafka 读取实时日志val ssc new StreamingContext(sparkConf, Seconds(5)) val lines: DStream[String] KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) ).map(_.value())这里的lines就是一个 DStream你对它做的每一次 transform 操作比如map、filter、reduceByKeySpark 都会把它翻译成对底层每个 RDD 的同样的操作。所以 DStream 本身不存储数据它只是描述了一个“每隔 5 秒要做什么”的计算图。DStream 的操作分为两类无状态转化操作和有状态转化操作。无状态操作很好理解每个批次之间互相独立算完就完有状态操作则需要跨批次累积数据典型的是updateStateByKey和reduceByKeyAndWindow。有状态操作跑起来之后每一个批次的结果都依赖以前批次的历史数据这就是为什么 Streaming 作业必须设置检查点不然状态一丢整个计算链就得从零开始。1.3 窗口与滑动时间必须遵守的倍数约束窗口操作是流计算里最常用的概念。假设业务要看过去 30 秒内每个商品的点击量每 10 秒刷新一次展示这时候就用到窗口。窗口长度是 30 秒滑动间隔是 10 秒。用 Spark Streaming 写出来是这样的val clicks: DStream[(String, Long)] lines .filter(_.contains(click)) .map(line (extractProductId(line), 1L)) val windowedCounts clicks.reduceByKeyAndWindow( (a: Long, b: Long) a b, (a: Long, b: Long) a - b, // 反向函数滑动窗口移出数据时做减法 Seconds(30), Seconds(10) )注意这里的反向函数(a, b) a - b很多人第一次写会漏掉。它的作用是新一批数据进入窗口时加上老一批数据滑出窗口时减去这样可以复用前一个窗口的计算结果避免每个窗口都重新算一遍全部数据。如果不提供反向函数Spark 只能把每个窗口视为独立批次重新计算数据量大时性能差距非常明显。窗口长度和滑动间隔有两个硬性约束窗口长度必须是批次间隔的整数倍滑动间隔也必须是批次间隔的整数倍。比如批次间隔是 5 秒窗口长度就不能设成 17 秒滑动间隔也不能设成 3 秒。这个约束的本质是 DStream 的数据粒度已经被批次间隔固定了窗口只能按整数个批次来拼。实际调参时窗口越长需要驻留内存的数据越多GC 压力越大滑动间隔越短计算频率越高CPU 开销越大。我的建议是先从“窗口时长 : 滑动间隔 3 : 1”这个比例起步比如 30 秒窗口配 10 秒滑动跑稳定后再尝试更激进的配置。2. 从 Kafka 到 Spark Streaming一个订单统计案例全流程落地2.1 环境准备与依赖别再踩版本坑真正做项目时第一步不是写代码而是把 Spark、Kafka 和 Scala 的版本对齐。Spark Streaming 接 Kafka 用的是专门的 connector它和 Kafka 的客户端版本强相关。以当前最常见的组合为例Spark 2.4.x 配 Kafka 0.10.x 连接器Scala 版本选 2.11 或 2.12。如果是 Spark 3.x连接器同样支持但因为 Spark Streaming 已经进入维护模式新项目我更推荐直接用 Spark 3.x 里的 Structured Streaming。Maven 依赖长这样dependency groupIdorg.apache.spark/groupId artifactIdspark-streaming_2.12/artifactId version2.4.8/version /dependency dependency groupIdorg.apache.spark/groupId artifactIdspark-streaming-kafka-0-10_2.12/artifactId version2.4.8/version /dependency注意那个spark-streaming-kafka-0-10的包名里面含0-10字样把很多不熟悉的人绕晕过。它表示的是 Kafka broker 的版本是 0.10 及以上不是 Spark 的版本。如果你的 Kafka 是 2.x 甚至 3.x用的依然是这个0-10连接器。集群这块如果你是自己搭测试环境Standalone 模式起步最简单spark-env.sh里配好 Java 和 SPARK_MASTER_HOST 就能跑。生产环境一般走 YARN提交命令通常是spark-submit \ --class com.example.OrderStreamingJob \ --master yarn \ --deploy-mode cluster \ --executor-memory 8g \ --num-executors 8 \ --executor-cores 4 \ order-job.jar我见过不少同事在这步栽跟头本地 IDEA 里跑得通一上集群就报ClassNotFoundException原因是依赖没打成 fat jar。建议直接用maven-shade-plugin把 spark-streaming-kafka-0-10 及其传递依赖打进去注意排除 Spark 自身的类。2.2 Direct 模式 vs Receiver 模式选错等于给自己挖坑Spark Streaming 接 Kafka 有两种方式Receiver 模式和 Direct 模式。我强烈建议新代码一律用 Direct 模式老代码如果还在用 Receiver能改就尽早改。Receiver 模式是老一代做法启动一个长驻的 Receiver 去 Kafka 拉数据数据先存进 Spark 的 block manager再交给批次计算。它的问题在于Receiver 挂了会丢数据需要靠 Write-Ahead Log 补但 WAL 又带来额外的磁盘开销而且因为数据经过了一层中转和 Kafka 分区的对应关系也丢了做精确一次语义非常困难。Direct 模式在 Kafka 0.10 连接器里是官方主推方式。它直接把 Kafka 的每个分区映射成 Spark RDD 的每个分区一个批次要处理哪些 offset在 Driver 端就能拿到不再需要 Receiver 中转。这样做有三个明显收益天然和 Kafka 分区对齐并行度可控不需要 WAL减少一环磁盘 IO通过手动提交 offset可以做到精确一次语义至少一次加幂等。下面这张表可以帮你做判断对比项Receiver 模式Direct 模式数据存储路径Kafka → Receiver → Block ManagerKafka → Spark Executor并行度受 Receiver 数量限制与 Kafka 分区数一一对应精确一次语义很难实现配合幂等可做到WAL 开销需要不需要适用建议存量老系统新系统一律选它2.3 完整代码消费、解析、窗口聚合、结果下沉我拿一个电商场景举例实时统计每分钟的订单总金额和订单量结果写入 MySQL供大屏展示。完整的 Direct 模式代码如下。import org.apache.kafka.clients.consumer.ConsumerRecord import org.apache.kafka.common.serialization.StringDeserializer import org.apache.spark.SparkConf import org.apache.spark.streaming.{Seconds, StreamingContext} import org.apache.spark.streaming.kafka010._ object OrderStreamingJob { def main(args: Array[String]): Unit { val conf new SparkConf() .setAppName(OrderStreamingJob) .set(spark.serializer, org.apache.spark.serializer.KryoSerializer) val ssc new StreamingContext(conf, Seconds(5)) // 生产环境检查点必须设置到 HDFS ssc.checkpoint(hdfs://path/to/checkpoint/order-streaming) val kafkaParams Map[String, Object]( bootstrap.servers - kafaka01:9092,kafka02:9092, key.deserializer - classOf[StringDeserializer], value.deserializer - classOf[StringDeserializer], group.id - order-streaming-group, enable.auto.commit - (false: java.lang.Boolean), auto.offset.reset - latest, spark.streaming.kafka.maxRatePerPartition - 2000 ) val topics Set(order-topic) val stream: InputDStream[ConsumerRecord[String, String]] KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) ) // 解析 JSON 并提取金额和单号 val orders stream.map(_.value()) .map(line { // 生产环境用 Jackson 或 Gson这里示意 val amount extractAmount(line) (1, amount) // 用固定 key 聚合全部订单 }) // 每 60 秒窗口每 10 秒滑动 val windowed orders.reduceByKeyAndWindow( (a: Double, b: Double) a b, (a: Double, b: Double) a - b, Seconds(60), Seconds(10) ) // 统计订单量每条记录算 1 单 val countStream orders.map { case (_, amount) (1, 1L) } val windowedCount countStream.reduceByKeyAndWindow( (a: Long, b: Long) a b, (a: Long, b: Long) a - b, Seconds(60), Seconds(10) ) windowed.join(windowedCount).foreachRDD { rdd rdd.foreachPartition { iter // 在 Executor 端建立 MySQL 连接一次批次写一批 val conn createDBConnection() iter.foreach { case (_, (amount, count)) upsertOrderSummary(conn, amount, count) } conn.close() } } ssc.start() ssc.awaitTermination() } }几个关键点要说明。第一enable.auto.commit必须设成false。Direct 模式下 offset 的提交由你决定如果自动提交一旦数据还没处理完进程就重启offset 会先于数据处理提前提交导致丢数据。这也是精确一次语义的第一个前提。第二spark.streaming.kafka.maxRatePerPartition是每分区最大读取速率这个参数在后面调优里会反复用到初次部署保守一点设个值能防止 Kafka 突发流量直接把作业打爆。第三落地到外部存储时不要在 Driver 端建连接然后在算子闭包里用它。连接要在foreachPartition里每个分区单独建这样才能保证连接跟随 Executor 分布不会形成单点瓶颈。第四也是最重要的仅仅靠 Streaming 框架你做到的是“至少一次”语义不是精确一次。如果下游写入 MySQL 时中途失败“至少一次”会导致同一批数据被重复写入。解决方案是让写入幂等比如给订单表加唯一键用INSERT ... ON DUPLICATE KEY UPDATE这样即便重复读到了同一批数据最终结果也不会翻倍。2.4 并行度与批次间隔跑起来之后的第一轮优化代码跑通之后第一个要调的就是并行度。Direct 模式下每个 Kafka 分区默认对应一个 RDD 分区一个分区只起一个 task。如果你的 Kafka topic 有 8 个分区就算 Executor 上有 16 个核真正并行的也只有 8 个 task剩下 8 个核白白闲着。这时可以在中间步骤加repartitionval repartitioned orders.repartition(16)但注意repartition是 shuffle 操作代价不低。如果目标分区数只是小幅增加可以用coalesce它尽量不 shuffle。我的习惯是先确认 Kafka 分区数再定 Executor 核数让两者尽量匹配或者让 Spark 分区数等于 Core 数的整数倍避免无谓 shuffle。批次间隔的选择也很有讲究。设得太短比如 1 秒Spark 每秒钟都要调度一次任务、产生一批新的 RDD调度开销占比会非常高真正算数据的时间反而少了。设得太长比如 60 秒延迟就大到没法看。比较靠谱的做法是先设 5 秒或 10 秒跑一批观察 Spark UI 上的“Batch Processing Time”如果每个批次实际处理时间不到间隔的一半再尝试缩短间隔如果处理时间已经超过 70% 的间隔说明系统接近瓶颈应该把间隔拉长而不是强行设小值硬扛。3. 内存模型、背压与调优让任务稳定跑下去3.1 流处理为什么会比批处理更容易 OOMSpark 的内存模型从 2.x 开始就不再区分静态的 execution 和 storage 区而是统一用 unified memory两个区域可以互相抢占。这是批处理的好消息对流处理却不一定是好消息。流处理的特殊性在于窗口数据天然需要“留在内存里”。比如窗口长度 60 秒批次间隔 5 秒那么任何时刻都有 12 个批次的数据在内存里等待计算。如果你的数据量是每批 10 万条窗口区就要同时维持 120 万条再加上计算结果、缓存、序列化临时对象内存消耗会迅速膨胀。我上面那个订单案例里如果用默认的MEMORY_ONLY存储级别RDD 的中间结果全部留在内存GC 会疯了一样地跑。这时的内存压力有两条出路调大 Executor 内存比如--executor-memory 8g起步用 Kryo 序列化替代 Java 序列化显著降低对象占用配置 Kryo 只需要两行conf.set(spark.serializer, org.apache.spark.serializer.KryoSerializer) conf.set(spark.kryo.registrationRequired, true)再显式注册你的业务类比如订单样例类。另外spark.memory.fraction这个参数决定了统一内存区占整个堆的比例默认是 0.6。流任务内存吃紧时我会把spark.memory.storageFraction调低比如从默认的 0.5 降到 0.2让更多内存留给执行计算避免窗口缓存无限挤压执行空间。这两个参数在不同版本的包里位置略有差异调之前先看当前版本的 Spark 配置文档。3.2 限流与背压让消费速度匹配处理速度流处理作业最怕什么最怕 Kafka 的流量突然翻倍消费速度跟不上生产速度处理不完的数据越积越多最后把 Executor 撑爆。解决这个问题Spark Streaming 提供了限流和背压两套机制。先看限流。spark.streaming.kafka.maxRatePerPartition直接限制每个分区每秒最多读取多少条消息。比如 Kafka 有 8 个分区设成 2000那么每个批次最多读2000 * 8 * 5秒 80000条。这个参数的好处是硬性可控坏处是静态的流量低了它不会自动放大流量高了你得手动改。再看背压。开启背压后系统会动态估算合适的消费速率自动调节每个批次拉取的数量不用你手工干预。配置如下conf.set(spark.streaming.backpressure.enabled, true) conf.set(spark.streaming.backpressure.initialRate, 1000)开启背压后maxRatePerPartition仍可作为速率上限兜底。系统会用 PID 控制器计算速率原则是宁可清空积压也不能让批处理时间无限膨胀。实际生产中我的做法是初始部署先手动限流观察每个批次的处理耗时和调度延迟等数据模式摸清了再开背压并保留一个上限。千万不要一开始就依赖背压它估出来的速率需要时间收敛期间可能先经历一段过载。这里有一个常见误区开了背压就觉得高枕无忧了。实际上背压只能控制拉取速率如果你下游的写入本身很慢比如 MySQL 批量插入吞吐太低背压会一直往下降速直到 Kafka 消费 lag 越来越大但问题其实出在 sink不是在源端。排坑顺序要遵循“先看下游再看上游”的原则。3.3 一次真实的内存溢出排查过程我以前排查过一个比较典型的 OOM 问题。现象是 Spark UI 里批处理时间从 8 秒逐步爬升到 20 秒GC 时间占比肉眼可见地增加最后几个 Executor 报java.lang.OutOfMemoryError: Java heap space。初步怀疑是窗口数据量太大于是打开 Executor 的 GC 日志发现 Full GC 频率大约是每 30 秒一次单次耗时 5 秒以上。同时看到 Spark UI 里的 Storage Memory 占用接近饱和说明 RDD 缓存把存储区占满了执行区抢不到内存任务只能反复等 GC。处理分了三步第一步器亡内存和序列化。把--executor-memory从 4g 调到 8g开启 Kryo 序列化注册了业务类。这一步执行后Full GC 频率降到了一分钟一次批处理时间回落到 12 秒左右。第二步限流。给maxRatePerPartition设了上限把每批数据量压到原来的 60%批处理时间降到 8 秒以内系统不再处于“边消费边积压”的状态。第三步优化窗口计算。把原来 60 秒窗口、每 5 秒滑动一次的配置改成 60 秒窗口、每 15 秒滑动一次并把输出目标从逐条写入改为每窗口批量写入减少连接开销。三步做完作业稳定运行了很长一段时间。这个案例想说明的是OOM 不会只有一个原因内存、速率、计算逻辑互相耦合调优也要一套组合拳。4. 高频故障、检查点与 API 演进踩坑记录与选型建议4.1 高频故障速查表下面这些故障是我在实际维护 Stream 作业时反复遇到的我整理成表方便你对照排查。症状可能原因解决办法批次处理时间持续大于批次间隔整体流量过载、窗口过大、下游写入慢提升分区并行度开限流或背压检查下游 sink重启后重复消费大量数据未手动提交 offset或提交时机不对改为手动提交确保业务处理成功后再提交Task not serializable算子里引用了 Driver 端的非序列化对象把连接、对象在算子内部创建或声明为 transient窗口数据错乱看起来像丢了窗口长度与批次间隔不是整数倍检查窗口参数确保是批间隔的整数倍GC 频繁吞吐骤降内存不足、未启用 Kryo、storageFraction 过高调大内存开 Kryo降 storageFractionKafka lag 持续上涨消费能力不足或背压速率估计偏保守增加执行资源或分区适当提高 maxRate 上限报RateController相关异常背压开启但未设置初始速率设置 backpressure.initialRateDriver 一直重启检查点目录不可用或权限不足确认 HDFS 路径存在、可写、命名空间正确表格只是触发点真正排障时一定要看 Spark UI 的 Streams 页面。那个页面会直接展示每个批次从开始接收数据到处理完成的所有时间占比是排查 Streaming 问题的第一入口。4.2 检查点与精确一次语义Offsets 必须自己管负载敬业的 Streaming 作业检查点不是可选项是必选项。它的作用是持久化三样东西流处理的操作图即 DStream 的计算链条、配置参数、以及已处理的 offset。设置方法简单上面代码里已经出现了ssc.checkpoint(hdfs://path/to/checkpoint)但你要非常清楚它的限制基于旧的检查点恢复的作业不能改变代码结构和计算逻辑。比如你原来只有map和filter恢复后改成map、filter再加上window操作图对不上任务会起不来。官方文档对此提过但很多人都忽略了等踩到坑才后悔。所以代码逻辑有变动时别直接用旧 checkpoint要么清掉重建要么就把它当成仅从最近 offset 恢复的手段配合外部存储维护的 offset 一起做。精确一次语义这里把话说透。Direct 模式 手动提交 offset 能保证“至少一次”处理意思是每条数据至少被处理一次可能被处理多次。要变成“精确一次”必须配合下游的幂等写入。上面订单案例里用唯一键 upsert就是幂等的一种实现。整个链路是三个环节源端可重放Kafka 天然满足、计算可重放Spark 天然满足、Sink 幂等自己实现。做好第三条才能对外宣称精确一次。另一种做法是“原子写入加事务”比如在同一个事务里写结果数据和提交 offset。但 Spark Streaming 因为分布式的特性跨系统事务很难做到全局原子实操中我还是推荐幂等方案简单、可靠、不容易出幺蛾子。4.3 新项目还用 Spark Streaming 吗Structured Streaming 的迁移思路写到这里必须要回答一个现实问题现在新项目还该不该学 Spark Streaming我的答案是作为概念必须要懂但新项目尽量用 Structured Streaming。Structured Streaming 是 Spark 2.0 引入的新一代流处理 API核心思想是流批统一把流当成一张无限追加的表用 DataFrame/Dataset 的 API 操作背后走 Spark SQL 引擎。它解决了 Spark Streaming 几个让人头疼的痛点事件时间和水位线是内置概念乱序数据有了标准化处理方式处理语义更清晰输出模式append、update、complete一目了然和 Kafka 的连接配置更简洁offset 管理更透明。而 Spark Streaming 的 DStream API 从 Spark 3.0 开始进入维护模式官方不再加新功能你后面用的每一年都在为旧技术还债。所以新项目选 Structured Streaming 是理性选择。迁移思路其实不难核心是把 DStream 的 transform 操作换成 DataFrame 的算子。数据源从KafkaUtils.createDirectStream变成spark.readStream.format(kafka)窗口操作从reduceByKeyAndWindow变成groupBy(window(col(ts), 60 seconds), col(key)).sum(amount)。业务逻辑不用推翻重来但代码确实要重写一遍。如果你手里正维护着一套跑了两三年的 Spark Streaming 作业而且业务稳定、没出过大事故我的建议是别急着重构。先把状态管理、checkpoint、监控体系梳理清楚新需求再用新 API 实现逐步过渡。线上系统的价值是稳定不是为了技术时髦去承担不必要的迁移风险。我自己在维护三套 Spark Streaming 作业的过程中最有感触的一点是框架怎么换流处理的几个核心问题——延迟与吞吐的取舍、状态管理、乱序数据、下游幂等——永远是逃不掉的。老老实实把这些问题摸透换什么 API 都能快速上手。最后给你留一个小技巧不管是 Spark Streaming 还是 Structured Streaming一定把每批次的处理时间、Kafka lag、GC 时间这三项指标接入你们的监控告警先于用户发现问题这是流作业运维里最重要的一条经验。
返回列表