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

Scala中如何向Kafka的commitAsync方法传递回调?附测试方案

Handling Kafka's commitAsync Callback in Scala & Writing Mockito Tests

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.
  • ArgumentCaptor lets us capture the callback passed to commitAsync, so we can manually trigger its onComplete method to test both success and failure paths.
  • verify(mockConsumer).commitAsync(...) ensures that our code actually calls the consumer's commitAsync method with the callback.

内容的提问来源于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 08:57:46