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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 09:14:43