【问题标题】:spring kafka looking for the latest available message in a topicspring kafka 在主题中寻找最新的可用消息
【发布时间】:2019-01-10 17:20:18
【问题描述】:

我有一个 Kafka 消费者,它订阅了以下主题 MY_TOPICMY_UNINTERESTED_TOPIC

在下面的场景中,我对第二个主题不感兴趣,但我不得不提及它,因为如果我使用 auto.offset.reset 之类的东西配置它,它可能会影响所有主题。

MY_TOPIC 主题上,我发布了不同类型的消息:MESSAGE_TYPE_AMESSAGE_TYPE_B。两条消息都是BaseKafkaMessage(自定义类)的实例,具有不同的属性。

现在我有兴趣找到 MESSAGE_TYPE_A 类型的最新消息。我该怎么做?

真正的场景是这样的:我在同一个主题上发布两种类型的消息。其中一个用于在每个对此主题和该消息感兴趣的消费者中准备一个本地缓存。如果消费者停止,当它重新加载时,它必须使用最新的MESSAGE_TYPE_A 重新初始化其缓存。 MESSAGE_TYPE_B 应该被忽略。我不想在 Kafka 上向数据提供者发送通知以再次发布数据,因为所有订阅者都会有很多不必要的工作要做。

我怎样才能获得这个?这可能吗?

我找到了https://docs.spring.io/spring-kafka/reference/htmlsingle/#seek,但我不确定这是否是我正在寻找的,或者是否有其他方法可以做到这一点。

【问题讨论】:

  • 但是每当消费者重启时,它会拉取新数据对吗?那为什么还要担心MESSAGE_TYPE_B,而auto.offset.reset这个属性如果是新的消费群会受到影响
  • @Deadpool 请看看 cricket_007 提供的答案以及我的评论
  • 您对消费者组 ID 有任何了解吗?我真的不明白消费者在使用 sama group id 重新启动时如何从一开始就进行轮询? @tzortzik
  • 欢迎使用帖子旁边的复选标记接受答案

标签: java apache-kafka kafka-consumer-api spring-kafka


【解决方案1】:

目前还不清楚这些消息采用什么格式,或者为什么您真的需要它们出现在同一个主题中。

例如,您可以use different Avro types。或者你必须try-catch 解析两个不同的字节数组(JSON 对象?)

或者您可以按类型对主题进行分区,这样一种类型的所有消息都按同一分区的顺序排列。

但是,没有进行索引查找以获取最近发送的消息的机制。要么从最新事件开始,然后收到下一条传入消息,要么从头开始,然后向上扫描,直到在下一个轮询循环中获得 0 条记录,这在理论上是“最近的”

类似 auto.offset.reset 的东西可能会影响所有主题

这只影响消费者感兴趣的主题,而不是全部。

【讨论】:

  • Avro 目前不是我的解决方案。我不想使用auto.offset.reset,因为我的消费者订阅了ORDER_MANAGEMENTPRODUCT_MANAGEMENT 之类的主题。在这种情况下,我们可能会重新创建之前创建的订单或产品。我只对带有类型消息的PRODUCT_MANAGEMENT 感兴趣(类型是消息上的一个字段)`PRODUCT_CACHE_UPDATE. The solution you propose is to search from the beginning until I have no messages and look for my latest PRODUCT_CACHE_UPDATE`?
  • 去重记录取决于您的消费者逻辑。是否更改 kafka 属性不会阻止它的发生。您可以使用 Kafka Streams 将您的一个主题 filter 根据类型分为 2 个主题。不过,再次“看”的行为需要进行主题扫描
猜你喜欢
  • 2021-05-12
  • 2015-12-07
  • 2019-06-08
  • 1970-01-01
  • 2022-07-28
  • 2020-01-07
  • 2020-02-08
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多