You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何从Spark Streaming获取committedOffsets、availableOffsets及代码获取偏移集合

如何从Spark Streaming中获取committedOffsets和availableOffsets?

日志示例

22/11/09 11:08:40 INFO MicroBatchExecution: Resuming at batch 206 with committed offsets {KafkaV2[Subscribe[test]]: 
{"test":{"0":3086,"1":3086,"2":3086,"3":3086,"4":3086,"5":3086,"6":3086,"7":3086,"8":3086,"9":3086,"10":3086,
"11":3086,"12":3086,"13":3086,"14":3086,"15":3086,"16":3086,"17":3086,"18":3086,"19":3086,"20":3086,"21":3086,"22":3086,"23":3086,
"24":3086,"25":3086,"26":3086,"27":3086,"28":3086,"29":3086,"30":3086,"31":3886,"32":3086,"33":3086,"34":3086,"35":3086,"36":3086,
"37":3086,"38":3086,"39":3086,"40":3086,"41":3086,"42":3086,"43":3086,"44":3086,"45":3086,"46":3086,"47":3086,"48":3086,"49":3086}}} 
and available offsets {KafkaV2[Subscribe[test]]: {"test":{"0":3105,"1":3105,"2":3105,"3":3105,"4":3105,
"5":3105,"6":3105,"7":3105,"8":3105,"9":3105,"10":3105,"11":3105,"12":3105,"13":3105,"14":3105,"15":3105,"16":3105,"17":3105,"18":3105,"19":3105,
"20":3105,"21":3105,"22":3105,"23":3105,"24":3105,"25":3105,"26":3105,"27":3105,"28":3105,"29":3105,"30":3105,"31":3910,"32":3105,"33":3105,"34":3105,
"35":3105,"36":3105,"37":3105,"38":3105,"39":3105,"40":3105,"41":3105,"42":3105,"43":3105,"44":3105,"45":3105,"46":3105,"47":3105,"48":3105,"49":3105}}}

上述偏移集合属于Spark Structured Streaming(日志中MicroBatchExecution为该框架的执行器)的Checkpoint偏移量,以下是几种获取方式:

1. 通过StreamingQuery的lastProgress获取

启动流查询后,直接调用lastProgress属性即可获取最近一批次的执行详情,其中包含已提交(committed)和可用(available)偏移量:

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.streaming.StreamingQuery

val spark = SparkSession.builder()
  .appName("OffsetFetch")
  .master("local[*]") // 生产环境移除该配置
  .getOrCreate()

// 构建Kafka流查询
val query: StreamingQuery = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "your-kafka-host:port")
  .option("subscribe", "test")
  .load()
  .writeStream
  .format("console")
  .option("checkpointLocation", "/path/to/your/checkpoint-dir")
  .start()

// 获取并打印偏移量
val progress = query.lastProgress
if (progress != null) {
  val sourceProgress = progress.sources.head
  println(s"已提交偏移量: ${sourceProgress.commitOffsets}")
  println(s"可用偏移量: ${sourceProgress.availableOffsets}")
}

// 保持查询运行,按需添加
query.awaitTermination()

2. 自定义StreamingQueryListener实时监听

如果需要实时监控每一批次的偏移量变化,可以自定义监听器:

import org.apache.spark.sql.streaming.StreamingQueryListener
import org.apache.spark.sql.streaming.StreamingQueryListener._

val offsetListener = new StreamingQueryListener {
  override def onQueryStarted(event: QueryStartedEvent): Unit = {}

  override def onQueryProgress(event: QueryProgressEvent): Unit = {
    val progress = event.progress
    val source = progress.sources.head
    println(s"批次ID: ${progress.batchId}")
    println(s"已提交偏移量: ${source.commitOffsets}")
    println(s"可用偏移量: ${source.availableOffsets}")
  }

  override def onQueryTerminated(event: QueryTerminatedEvent): Unit = {}
}

// 注册监听器
spark.streams.addListener(offsetListener)

// 启动查询(同方法1的query代码)
query.awaitTermination()

3. 直接解析Checkpoint目录的偏移量文件

Checkpoint目录的offsets子目录下会保存每一批次的已提交偏移量文件(文件名对应批次ID),可以直接读取解析:

import scala.io.Source
import org.json4s._
import org.json4s.jackson.JsonMethods._

val checkpointDir = "/path/to/your/checkpoint-dir"
val offsetsDir = new java.io.File(s"$checkpointDir/offsets")

// 获取最新的偏移量文件
val latestOffsetFile = offsetsDir.listFiles().sortBy(_.getName.toLong).last
val offsetJson = Source.fromFile(latestOffsetFile).mkString

// 解析JSON格式的偏移量
implicit val formats = DefaultFormats
val committedOffsets = parse(offsetJson).extract[Map[String, Any]]
println(s"解析后的已提交偏移量: $committedOffsets")

注意事项

  • 若使用的是旧版DStream API(传统Spark Streaming),获取偏移量的方式不同,可通过KafkaUtils.createDirectStream的OffsetRange对象或者Checkpoint中的元数据恢复。
  • 可用偏移量本质是Kafka集群当前的最新偏移,也可以通过Kafka客户端API直接查询。

内容的提问来源于stack exchange,提问作者qtcat

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.13 03:05:34