如何在Scala中实现Kafka的ConsumerRebalanceListener接口(订阅主题触发)
Implementing Kafka's ConsumerRebalanceListener in Scala for Topic Subscriptions
Hey there! Implementing Kafka's ConsumerRebalanceListener in Scala is totally straightforward, thanks to Scala's seamless interoperability with Java. Whether you want a one-off listener for a single subscription or a reusable class, here's how to do it:
1. Anonymous Listener (Most Common for Direct Subscription)
When you subscribe to topics, you can directly define an anonymous implementation of ConsumerRebalanceListener right in the subscribe call. This is perfect for simple, one-off logic:
import org.apache.kafka.clients.consumer.{ConsumerRebalanceListener, KafkaConsumer, TopicPartition} import java.util import scala.jdk.CollectionConverters._ // Initialize your KafkaConsumer first (props should include bootstrap servers, group id, etc.) val consumer = new KafkaConsumer[String, String](props) // Subscribe to topics with the rebalance listener consumer.subscribe( util.Arrays.asList("your-topic-1", "your-topic-2"), new ConsumerRebalanceListener { override def onPartitionsRevoked(partitions: util.Collection[TopicPartition]): Unit = { // Logic to run BEFORE partitions are revoked (critical for offset commits!) val scalaPartitions = partitions.asScala.map(tp => s"${tp.topic()}:${tp.partition()}") println(s"Partitions being revoked: ${scalaPartitions.mkString(", ")}") // Example: Commit offsets synchronously to avoid duplicate messages consumer.commitSync() } override def onPartitionsAssigned(partitions: util.Collection[TopicPartition]): Unit = { // Logic to run AFTER new partitions are assigned val scalaPartitions = partitions.asScala.map(tp => s"${tp.topic()}:${tp.partition()}") println(s"New partitions assigned: ${scalaPartitions.mkString(", ")}") // Example: Reset offset to the beginning of the assigned partitions // partitions.asScala.foreach(tp => consumer.seek(tp, 0L)) } } )
2. Reusable Custom Listener Class
If you need to reuse the rebalance logic across multiple consumers, create a dedicated class that implements ConsumerRebalanceListener:
import org.apache.kafka.clients.consumer.{ConsumerRebalanceListener, KafkaConsumer, TopicPartition} import java.util import scala.jdk.CollectionConverters._ class CustomRebalanceHandler(consumer: KafkaConsumer[String, String]) extends ConsumerRebalanceListener { override def onPartitionsRevoked(partitions: util.Collection[TopicPartition]): Unit = { // Custom revocation logic - e.g., async offset commit with error handling println(s"Revoking ${partitions.size()} partitions") consumer.commitAsync((offsets, exception) => { if (exception != null) { println(s"Failed to commit offsets during rebalance: ${exception.getMessage}") } }) } override def onPartitionsAssigned(partitions: util.Collection[TopicPartition]): Unit = { // Custom assignment logic - e.g., load offsets from a database val partitionDetails = partitions.asScala.map(tp => s"${tp.topic()}[${tp.partition()}]").mkString(", ") println(s"Assigned partitions: $partitionDetails") // Example: Fetch offset from external storage and set it // partitions.asScala.foreach(tp => { // val storedOffset = fetchOffsetFromDB(tp.topic(), tp.partition()) // consumer.seek(tp, storedOffset) // }) } // Helper method example (replace with your actual logic) private def fetchOffsetFromDB(topic: String, partition: Int): Long = { // Your database lookup logic here 0L } } // Usage: val consumer = new KafkaConsumer[String, String](props) val rebalanceHandler = new CustomRebalanceHandler(consumer) consumer.subscribe(util.Arrays.asList("your-topic"), rebalanceHandler)
Key Notes
- Java/Scala Collection Interop: Use
scala.jdk.CollectionConverters._to convert Javautil.Collectionto Scala collections (likeSeq), making it easier to work with using Scala's functional methods. - Critical Logic:
onPartitionsRevokedis where you should commit offsets to prevent duplicate messages after a rebalance.onPartitionsAssignedis where you can adjust offsets (e.g., reset to a specific point or load from external storage) before starting to consume. - Thread Safety: Remember that Kafka consumer methods are not thread-safe, so avoid calling consumer operations from outside the listener methods unless you handle synchronization properly.
内容的提问来源于stack exchange,提问作者cikavladimir
相关产品推荐
相关产品推荐

