【发布时间】:2020-11-16 14:57:53
【问题描述】:
我有一个使用 DLQ 方法的项目,在 @KafkaListener 中出现任何异常时,错误将被发送到结构为 error-<topic>-<consumergroup> 的主题。当我们想手动重试来自这个 DLQ 的 kafka 消息时,我们会将其生成到具有相似结构的主题 retry-<topic>-<consumergroup>。因此,例如,如果我们有一个主要主题 foo,由消费者组 bar 消费,我们将拥有以下主题:foo、error-foo-bar 和 retry-foo-bar。
这样,我们可以为特定的消费者组重试消息。消息处理程序从主主题和重试主题中读取。
由于我们在项目中有多个监听器都使用同一个容器,所以我们将主题(main、error 和 retry)放在应用程序属性中并配置DeadLetterPublishingRecoverer 以找到对应的DLQ(注意,这个是 Kotlin):
val recoverer = DeadLetterPublishingRecoverer(template) { record, _ ->
val dlq = kafkaProperties.topics.values.firstOrNull { it.main == record.topic() || it.retry == record.topic() }!!.dlq
return@DeadLetterPublishingRecoverer TopicPartition(dlq, -1)
}
这一切都很好。但是,现在我想添加另一个 @KafkaListener,它与另一个听众一样在 same topic 上收听。因此我需要给这个监听器一个单独的消费者组。
但是,如果我使用与上述相同的逻辑来查找相应的 DLQ 主题,则两个侦听器中的一个将使用另一个侦听器的 DLQ,因为它不是在这里查看消费者组。
所以我的问题是:在 Spring Kafka 中,有没有办法找到发生错误的消费者组(又名groupId)?我查看了文档和源代码,但我自己找不到。任何帮助表示赞赏。
【问题讨论】: