【发布时间】:2016-12-15 12:36:20
【问题描述】:
我需要对 kafka 应用程序执行单元测试,避免使用第三方库。
我现在的问题是我想清除测试之间的所有主题,但我不知道如何。
这是我的临时解决方案:提交每次测试后产生的每条消息,并将所有测试消费者放在同一个消费者组中。
override protected def afterEach():Unit={
val cleanerConsumer= newConsumer(Seq.empty)
val topics=cleanerConsumer.listTopics()
println("pulisco")
cleanerConsumer.subscribe(topics.keySet())
cleanerConsumer.poll(100)
cleanerConsumer.commitSync()
cleanerConsumer.close()
}
这不起作用,我不知道为什么。
例如,当我在测试中创建一个新的消费者时,messages 包含上一个测试中产生的消息。
val consumerProbe = newConsumer(SMSGatewayTopic)
val messages = consumerProbe.poll(1000)
我该如何解决这个问题?
【问题讨论】:
-
Kafka 是一个持久的消息存储,它让消费者决定开始消费的偏移量。你只需要记住你最后一次删除的偏移量,然后开始消费。
-
为什么这个问题用 Java 标记?结果,它出现在非常特定于 java 的搜索中。不是很有帮助。
标签: java scala unit-testing apache-kafka