能否提前配置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
相关产品推荐
相关产品推荐

