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

资讯详情

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

SparkStreaming 之接收数据原理剖析

SparkStreaming 之接收数据原理剖析 摘要做实时计算最先要搞明白的是数据怎么进来的。这篇讲清 Spark Streaming 的 Receiver 接收链路——数据从 Kafka/Socket 进来经 BlockGenerator 切成 Block 落进 BlockManager再按 batch 打包成 RDD顺带说清楚 batchInterval 和 blockInterval 这两个最容易混的参数以及 WAL、反压和 Receiver/Direct 两种模式的取舍。关键词Spark Streaming, Receiver, BlockGenerator, blockInterval, WAL, 反压, Direct 模式一、数据是怎么流进来的Spark Streaming 用的是微批处理——把实时流切成一个个小批次batch每个批次当成一个 RDD 来处理。所以接收数据的本质就是把持续到达的数据按时间窗口切块、攒批再交给后面的计算。一条数据的完整旅程是这样的数据源(Kafka/Socket) → Receiver → BlockGenerator → Block → BlockManager → WAL ↓ 每个 batch 的 Block 集合 → RDD → Job两个关键角色Receiver跑在 Executor 上的长驻进程负责从数据源持续拉取数据。它自己不吃数据只管搬。ReceiverTracker跑在 Driver 端管理所有 Receiver 的启停和元数据记录哪些 Block 属于哪个 batch。二、batchInterval 和 blockInterval别搞混这是理解 Spark Streaming 接收机制最重要的一对参数名字像作用完全不同。valsscnewStreamingContext(conf,Seconds(2))// batchInterval 2s// blockInterval 通过配置项设置默认 200msconf.set(spark.streaming.blockInterval,200ms)batchInterval批间隔每个 RDD 覆盖的时间窗口决定处理节奏。设 2 秒就是每 2 秒生成一个 RDD。blockInterval块间隔Receiver 把数据切成 Block 的粒度。默认 200ms。两者的关系决定了并行度一个 Block 对应 RDD 的一个分区。batchInterval 2 秒、blockInterval 200ms那么一个 batch 里有2000 / 200 10个 Block也就是这个 RDD 有 10 个分区最多 10 个并发任务。这里有个容易踩的坑如果单个 Receiver 的吞吐很高比如单 topic 单分区每秒几万条默认 200ms 的 blockInterval 会让每个 Block 塞进大量数据Block 撑爆内存直接 OOM。这时要么调大spark.streaming.blockInterval让 Block 更多更小要么干脆加 Receiver 数量把吞吐分散开。三、WAL数据不丢的代价Receiver 收到数据后默认是存在内存里的。一旦 Executor 挂了这些还没处理的数据就没了。要保证不丢得开 WALWrite Ahead Log预写日志conf.set(spark.streaming.receiver.writeAheadLog.enable,true)开了之后数据会先写进 WAL落盘再进 BlockGenerator。Receiver 挂了重启可以从 WAL 把没处理完的数据读回来。代价也很直接每次写盘都有 IO 开销吞吐会掉一截。所以 WAL 不是免费的是数据不丢和性能之间的取舍。这也是后面 Direct 模式要解决的点。四、反压接收太快处理不过来怎么办生产里很常见数据高峰时Receiver 拉取速度远超下游处理速度Block 越积越多最后 OOM。反压Backpressure就是解决这个的conf.set(spark.streaming.backpressure.enabled,true)原理是 Spark 用一个 PID 控制器根据上一批的实际处理时长动态算出下一批应该拉多少数据。处理慢了就少拉一点让上下游回到平衡。注意反压只对 Receiver 模式有效且开启后拉取速率是动态调整的测试时要注意它可能掩盖掉处理逻辑本身的性能问题。五、Receiver 模式 vs Direct 模式这是 Kafka 集成的两条路线也是理解 Spark Streaming 演进的关键。Receiver 模式旧valkafkaStreamKafkaUtils.createStream(ssc,zkQuorum,group,topicMap)// 长驻 Receiveroffset 存 ZKReceiver 用 Kafka 的 High-Level Consumer 拉数据offset 交给 ZooKeeper 管理。问题在于offset 的提交和实际处理是分离的——处理失败了但 offset 已经提交数据就丢了开了 WAL 变成 at-least-once又可能重复消费。而且每个 Receiver 要占一个核资源开销也大。Direct 模式新推荐valkafkaStreamKafkaUtils.createDirectStream[String,String](ssc,PreferConsistent,Subscribe[String,String](topics,kafkaParams))Direct 模式没有 Receiver每个 batch 直接拿着 Kafka 分区的 offset 范围去拉数据offset 由 Spark 自己管理、存进 checkpoint而且 offset 的更新和 Job 的成功与否绑定——Job 没成功就不提交 offset天然得到 exactly-once。一句话判断新项目无脑用 Direct 模式。Receiver 模式基本只剩历史兼容意义到 Spark 2.x 的结构化流Structured Streaming里底层已经统一成类 Direct 的拉取方式了。六、总结接收数据的核心链路是 Receiver → BlockGenerator → Block → BlockManager一个 Block 就是一个分区。batchInterval 决定处理节奏blockInterval 决定分区粒度后者默认 200ms高吞吐场景要留意 Block 过大导致的 OOM。WAL 是数据不丢和性能的取舍Direct 模式用 offset 自管理绕开了这个两难。反压靠 PID 控制器动态限速是应对数据高峰的保命配置。作者大数据技术实践者博客blog.starzy.cnGitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践
返回列表