如何从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
相关产品推荐
相关产品推荐

