本节摘要:以 Kafka 为代表的消息队列通过流式连接器接入 Spark,第 3 章的微批引擎在源侧多了一个新变量——消费位移。本节讲分区与任务的对应关系、位移提交时机与两种失效语义,以及吞吐与延迟的配平参数。
文件和数据库里的数据是静止的,队列里的数据有"当前位置"。这一节把第 3 章的微批模型推进到真实最常见的源头:Driver 每个微批开始时向 Kafka 询问各分区的新消息量,规划出与分区对应的任务,批结束后提交位移。位移这个变量,决定了消息会不会丢、会不会重。
kafka_src = spark.readStream.format("kafka") \ .option("kafka.bootstrap.servers", "broker1:9092,broker2:9092") \ .option("subscribe", "user_events") \ .option("startingOffsets", "latest") \ .option("maxOffsetsPerTrigger", "500000") \ .load() parsed = kafka_src.selectExpr( "CAST(value AS STRING) AS json_str", "topic", "partition", "offset")
每个微批的规划发生在 Driver:查询各 topic 分区的最新位移,减去上次提交位,得到新增消息量,再被 maxOffsetsPerTrigger 封顶——这个参数是吞吐闸门,它把每个批的输入规模压在可控范围,代价是消息洪峰时延迟抬升。每个 Kafka 分区映射为一个输入任务,分区数因此成为源侧并行度的地板:想提高消费并行度,先加 Kafka 分区。
批的三种结局对应三种位移命运:
| 批的结局 | 位移动作 | 重启后果 |
|---|---|---|
| 成功完成检查点 | 位移随检查点持久化 | 从上次位置继续,不丢不重 |
| 失败且恢复 | 回滚到上个检查点 | 检查点之后的消息重消费 |
| 无检查点且依赖自动提交 | 时间片到了就提交 | 位移先于处理完成,可能丢消息 |
结论直白:Structured Streaming 的可靠消费押在检查点上,自动位移提交在生产流作业里基本不该用。恰好一次输出要求"位移与结果原子绑定"——文件接收器把位移写进每批目录的提交记录,Kafka 接收器用事务把结果消息与位移绑在一个事务里,各有一套实现,但原理同源。
# 检查点目录:位移、状态、Watermark 全靠它恢复 out = parsed.writeStream \ .format("parquet") \ .option("path", "hdfs://nn:9000/stream/events") \ .option("checkpointLocation", "hdfs://nn:9000/ckpt/events") \ .trigger(processingTime="30 seconds") \ .start()
先加分区。输入任务数被 Kafka 分区数锁死,分区不加、Executor 再多也只是空转排队;分区扩了之后若单批处理时间仍贴着间隔,才是加 Executor 的时机。顺序颠倒会买到一堆闲着的核,账面上集群翻倍、吞吐纹丝不动。
背景:事件流写 Parquet 仓,检查点 30 秒一批,下游按小时聚合。操作:某夜作业因 Executor lost 重启两次。结果:次日聚合发现 11:42 前后两分钟的事件出现双份。定位:两次失败都发生在批处理中途,恢复回滚到上个检查点,最后一批消息被重消费——语义上是"至少一次",写文件接收器没有按批目录去重的下游设计。处置:下游按 batchId 加事件 id 做幂等合并;根治则给关键告警配自动拉起。解读:失效语义不是 bug 而是设计属性,工程上的正确姿势是让下游幂等,而不是幻想流处理恰好一次落仓。
# 三个最常动的旋钮 .option("maxOffsetsPerTrigger", "800000") # 批的输入上限 .option("fetch.max.bytes", "52428800") # 单次拉取字节 .trigger(processingTime="15 seconds") # 批间隔
配平逻辑:批间隔应大于批处理时间,留 20% 余量防积压累积;积压持续增长时优先加分区与 Executor,其次才调大批上限——单批太大会把第 3 章讲过的调度延迟放大成肉眼可见的卡顿。反方向,追求低延迟把间隔压到秒级,要警惕任务调度固定开销占比上升,微批模型的物理下限就在这里。
⚠️ 常见坑:Kafka 分区重分配后,旧检查点里的位移元数据可能指向已迁移的分区leader,消费端报 leader not available。巡检顺序:先看 broker 侧再认,再清 Kafka 参数缓存,最后才考虑重置位移。
💡 关键直觉:把"分区—任务—位移"想成三位一体。分区定并行度,任务按分区领数据,位移记每个分区吃到哪。三者任何一处对不上号,丢重问题就在那里发生。
最后一节看两个改变引擎边界的新构件:Spark Connect 重画客户端与集群的分界线,Delta Lake 给存储层补上事务。