【问题标题】:Is there a way to find consumer group in DeadLetterPublishingRecoverer?有没有办法在 DeadLetterPublishingRecoverer 中找到消费者组?
【发布时间】:2020-11-16 14:57:53
【问题描述】:

我有一个使用 DLQ 方法的项目,在 @KafkaListener 中出现任何异常时,错误将被发送到结构为 error-<topic>-<consumergroup> 的主题。当我们想手动重试来自这个 DLQ 的 kafka 消息时,我们会将其生成到具有相似结构的主题 retry-<topic>-<consumergroup>。因此,例如,如果我们有一个主要主题 foo,由消费者组 bar 消费,我们将拥有以下主题:fooerror-foo-barretry-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)?我查看了文档和源代码,但我自己找不到。任何帮助表示赞赏。

【问题讨论】:

    标签: apache-kafka spring-kafka


    【解决方案1】:

    消费者的group.id 可以通过调用KafkaUtils.getGroupId() 获得(容器线程启动时它存储在ThreadLocal 中)。

    【讨论】:

    • 嗨,Gary,感谢您的回答? 有没有其他选择或者这是首选的解决方案?似乎我必须弄清楚我是否可以在我现在拥有的 DLPR 配置中使用它(目前生活在 Bean 定义中)。
    • 您不需要“弄清楚我是否可以使用它”。 DLPR 始终在消费者线程上调用,因此在该线程上调用目标解析器时始终有效。
    猜你喜欢
    • 1970-01-01
    • 2020-11-10
    • 1970-01-01
    • 2020-11-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多