Scala中如何向Kafka的commitAsync方法传递回调?附测试方案
Hey there! Let's break down two key parts of your question: passing an OffsetCommitCallback to Kafka's Java API from Scala, and writing a corresponding test using MockitoSugar.
1. Passing OffsetCommitCallback in Scala
Kafka's OffsetCommitCallback is a Java functional interface (it has just one abstract method: onComplete). In Scala, you have two clean ways to implement and pass it:
Option 1: Anonymous Inner Class (Compatible with All Scala Versions)
If you're on Scala 2.11 or earlier, or prefer explicit code, use an anonymous class to implement the interface:
import org.apache.kafka.clients.consumer.{KafkaConsumer, OffsetCommitCallback, OffsetAndMetadata} import org.apache.kafka.common.TopicPartition def commitAsync(consumer: KafkaConsumer[String, String]): Unit = { consumer.commitAsync(new OffsetCommitCallback { override def onComplete( offsets: java.util.Map[TopicPartition, OffsetAndMetadata], exception: Exception ): Unit = { if (exception != null) { // Handle commit failure (log, retry, alert, etc.) println(s"Offset commit failed: ${exception.getMessage}") } else { // Handle successful commit println("Offsets committed successfully!") } } }) }
Option 2: Scala Lambda (Scala 2.12+)
Scala 2.12 and later support SAM (Single Abstract Method) conversion, which lets you use a lambda directly instead of an anonymous class. This keeps your code concise:
import org.apache.kafka.clients.consumer.{KafkaConsumer, OffsetAndMetadata} import org.apache.kafka.common.TopicPartition def commitAsync(consumer: KafkaConsumer[String, String]): Unit = { consumer.commitAsync((offsets: java.util.Map[TopicPartition, OffsetAndMetadata], exception: Exception) => { if (exception != null) { println(s"Offset commit failed: ${exception.getMessage}") } else { println("Offsets committed successfully!") } }) }
Pro tip: The compiler can often infer parameter types automatically, so you can simplify further:
consumer.commitAsync((offsets, exception) => { // Same logic as above })
2. Writing Tests with MockitoSugar
Let's assume your commit logic is wrapped in a class (like KafkaOffsetHandler) for testability. Here's how to write a test that verifies the commitAsync method is called correctly, and validates the callback behavior.
First, make sure your build.sbt includes the necessary test dependencies:
libraryDependencies ++= Seq( "org.mockito" %% "mockito-scala" % "1.17.12" % Test, "org.scalatest" %% "scalatest" % "3.2.15" % Test )
Now, the test class:
import org.apache.kafka.clients.consumer.{KafkaConsumer, OffsetCommitCallback, OffsetAndMetadata} import org.apache.kafka.common.TopicPartition import org.mockito.ArgumentCaptor import org.mockito.MockitoSugar._ import org.scalatest.BeforeAndAfterEach import org.scalatest.flatspec.AnyFlatSpec class KafkaOffsetHandlerSpec extends AnyFlatSpec with BeforeAndAfterEach with MockitoSugar { private var mockConsumer: KafkaConsumer[String, String] = _ private var offsetHandler: KafkaOffsetHandler = _ override def beforeEach(): Unit = { // Initialize mocks before each test mockConsumer = mock[KafkaConsumer[String, String]] offsetHandler = new KafkaOffsetHandler(mockConsumer) } "KafkaOffsetHandler.commitAsync" should "invoke consumer.commitAsync with a valid callback" in { // Call the method we want to test offsetHandler.commitAsync() // Capture the callback passed to commitAsync val callbackCaptor: ArgumentCaptor[OffsetCommitCallback] = ArgumentCaptor.forClass(classOf[OffsetCommitCallback]) verify(mockConsumer).commitAsync(callbackCaptor.capture()) // Test the callback's success scenario val capturedCallback = callbackCaptor.getValue val mockOffsets = mock[java.util.Map[TopicPartition, OffsetAndMetadata]] capturedCallback.onComplete(mockOffsets, null) // Add assertions here for your success logic (e.g., verify logs, service calls) // Test the callback's failure scenario val mockException = new RuntimeException("Commit failed unexpectedly") capturedCallback.onComplete(mockOffsets, mockException) // Add assertions here for your failure logic (e.g., verify error handling) } } // Example class containing our commit logic class KafkaOffsetHandler(consumer: KafkaConsumer[String, String]) { def commitAsync(): Unit = { consumer.commitAsync((offsets, exception) => { if (exception != null) { // Your actual failure handling logic } else { // Your actual success handling logic } }) } }
Key Test Details:
- We use
mock[KafkaConsumer[String, String]]to create a mock Kafka consumer, so we don't need a real Kafka cluster for testing. ArgumentCaptorlets us capture the callback passed tocommitAsync, so we can manually trigger itsonCompletemethod to test both success and failure paths.verify(mockConsumer).commitAsync(...)ensures that our code actually calls the consumer'scommitAsyncmethod with the callback.
内容的提问来源于stack exchange,提问作者Joe

