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

能否提前配置Kafka MockProducer处理send()调用?

能否提前配置Kafka MockProducer处理send()调用?

是的,默认的Kafka MockProducer不支持提前配置后续send()调用的处理逻辑。它的completeNext()和errorNext()方法只能作用于已经发起但尚未完成的send()请求——如果在send()调用前调用这些方法,不会产生任何效果,因为此时还没有待处理的Future需要完成。

你的测试挂起正是因为producer.send(...).get()会阻塞线程等待Future完成,但MockProducer默认不会自动完成这些Future,必须显式调用completeNext()来触发完成。而由于get()是阻塞的,测试主线程无法在send()调用后执行completeNext(),最终导致死锁。


可行解决方案

方案1:用Mockk直接模拟Producer接口

如果你习惯Mockito/Mockk的提前Stub方式,可以直接用Mockk模拟Kafka的Producer接口,跳过官方MockProducer。这样就能提前指定send()方法返回已完成的Future:

"Some test of random function" {
  val mockProducer = mockk<Producer<String, String>>()
  
  // 提前指定send()调用返回已完成的Future
  every { mockProducer.send(any()) } returns CompletableFuture.completedFuture(
    RecordMetadata(TopicPartition("test-topic", 0), 0, 0, 0, 0, 0)
  )

  val randomClass = RandomClass(producer = mockProducer)
  
  randomClass.randomFunction()    

  // 验证调用次数
  verify(exactly = 1) { mockProducer.send(any()) }
}

方案2:封装MockProducer实现预配置逻辑

如果必须使用官方MockProducer,可以自行封装一层,实现提前队列化完成/错误操作的逻辑:

class PreConfigurableMockProducer<K, V>(
    autoComplete: Boolean,
    keySerializer: Serializer<K>,
    valueSerializer: Serializer<V>
) : MockProducer<K, V>(autoComplete, keySerializer, valueSerializer) {

    private val pendingActions = ArrayDeque<() -> Unit>()

    // 提前加入"完成下一次send"的操作
    fun enqueueCompleteNext() {
        pendingActions.add { super.completeNext() }
    }

    // 提前加入"让下一次send报错"的操作
    fun enqueueErrorNext(exception: Exception) {
        pendingActions.add { super.errorNext(exception) }
    }

    override fun send(record: ProducerRecord<K, V>): Future<RecordMetadata> {
        val future = super.send(record)
        // 发送后立即执行预配置的操作
        pendingActions.poll()?.invoke()
        return future
    }
}

测试中使用这个封装类:

"Some test of random function" {
  val mockProducer = PreConfigurableMockProducer(
    false,
    StringSerializer(),
    StringSerializer()
  )

  val randomClass = RandomClass(producer = mockProducer)
  
  // 提前配置下一次send()完成
  mockProducer.enqueueCompleteNext()

  randomClass.randomFunction()    

  mockProducer.history().size shouldBe 1
}

方案3:异步线程触发完成操作

虽然不够优雅,但可以通过异步线程在send()调用后自动触发完成:

"Some test of random function" {
  val mockProducer = MockProducer(
    false,
    StringSerializer(),
    StringSerializer()
  )

  val randomClass = RandomClass(producer = mockProducer)
  
  // 启动异步线程等待send()发起后完成Future
  thread {
    while (mockProducer.history().isEmpty()) {
        Thread.sleep(10)
    }
    mockProducer.completeNext()
  }

  randomClass.randomFunction()    

  mockProducer.history().size shouldBe 1
}

总结:官方MockProducer的设计确实是针对已发起的请求进行处理,没有提供提前配置的API。如果你偏好Mockk/Mockito的风格,直接模拟Producer接口是最简洁的方案;如果必须使用官方MockProducer,封装一层实现预配置逻辑是更可靠的选择。

内容的提问来源于stack exchange,提问作者Stene

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 06:33:10