如何在Scala中实现Kafka相关Java接口及ConsumerRebalanceListener
嘿,我来帮你梳理下Scala里实现Kafka Java接口的方法,特别是你提到的ConsumerRebalanceListener~
在Scala中实现Kafka Java接口的核心思路
Scala和Java是完全兼容的,实现Java接口的方式非常直接,主要有两种风格:一是创建独立类实现接口,二是用匿名类(更贴合Scala的简洁写法)。下面重点针对你关心的重平衡监听器展开说明。
1. 实现ConsumerRebalanceListener的Scala类
对应你给出的Java示例SaveOffsetsOnRebalance,Scala里可以这样定义类:
import org.apache.kafka.clients.consumer.{ConsumerRebalanceListener, KafkaConsumer, OffsetAndMetadata} import org.apache.kafka.common.TopicPartition import scala.jdk.CollectionConverters._ // 传入消费者实例,方便在重平衡时操作偏移量 class SaveOffsetsOnRebalance(consumer: KafkaConsumer[String, String]) extends ConsumerRebalanceListener { // 重平衡开始前触发:通常在这里提交当前已处理的偏移量,避免重复消费 override def onPartitionsRevoked(partitions: java.util.Collection[TopicPartition]): Unit = { println(s"即将失去分区:${partitions.asScala.map(tp => s"${tp.topic()}-${tp.partition()}").mkString(",")}") // 同步提交偏移量,确保提交完成后再失去分区 consumer.commitSync() } // 重平衡完成后触发:可以在这里初始化新分区的消费位置 override def onPartitionsAssigned(partitions: java.util.Collection[TopicPartition]): Unit = { println(s"分配到新分区:${partitions.asScala.map(tp => s"${tp.topic()}-${tp.partition()}").mkString(",")}") // 示例:将新分区的消费位置重置到最开始 partitions.asScala.foreach(tp => consumer.seekToBeginning(List(tp).asJava)) } }
这里用到了scala.jdk.CollectionConverters._工具类,把Java集合转成Scala集合,用asScala/asJava实现互转,代码风格更贴合Scala习惯。
2. 订阅主题时使用该监听器的完整示例
下面是一个包含消费者初始化、订阅主题、使用重平衡监听器的完整代码:
import org.apache.kafka.clients.consumer.{ConsumerConfig, KafkaConsumer} import scala.jdk.CollectionConverters._ import java.util.Properties import java.time.Duration object KafkaConsumerWithRebalanceExample { def main(args: Array[String]): Unit = { // 配置消费者基础属性 val props = new Properties() props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092") props.put(ConsumerConfig.GROUP_ID_CONFIG, "scala-test-group") props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer") props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer") props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest") // 创建消费者实例 val consumer = new KafkaConsumer[String, String](props) // 实例化自定义重平衡监听器 val rebalanceListener = new SaveOffsetsOnRebalance(consumer) // 订阅主题并绑定监听器 consumer.subscribe(List("test-topic-1", "test-topic-2").asJava, rebalanceListener) try { // 持续消费消息 while (true) { val records = consumer.poll(Duration.ofMillis(100)) records.asScala.foreach(record => { println(s"收到消息:key=${record.key()}, value=${record.value()}, 分区=${record.partition()}") // 业务处理完成后,异步提交偏移量(性能更优) consumer.commitAsync() }) } } finally { // 关闭消费者前同步提交最终偏移量 consumer.commitSync() consumer.close() } } }
3. 简洁版:用匿名类实现监听器
如果你的重平衡逻辑比较简单,不想单独写类,可以直接用匿名类的方式实现:
consumer.subscribe(List("test-topic").asJava, new ConsumerRebalanceListener { override def onPartitionsRevoked(partitions: java.util.Collection[TopicPartition]): Unit = { println("重平衡前提交偏移量") consumer.commitSync() } override def onPartitionsAssigned(partitions: java.util.Collection[TopicPartition]): Unit = { println("重平衡后初始化新分区") } })
内容的提问来源于stack exchange,提问作者Joe
相关产品推荐
相关产品推荐

