【问题标题】:Unique messages from a kafka topic with in a time interval某个时间间隔内来自 kafka 主题的唯一消息
【发布时间】:2020-12-21 02:42:12
【问题描述】:

我有一个 kafka 主题,由不同的生产者产生 1500 条消息/秒,每条消息都有两个固定键 RID 和日期,(每条消息还有其他键不同)

有没有办法在主题中引入 1 分钟的延迟,并在 1 分钟的窗口中只消费唯一的消息。

示例 - 在一分钟内可能有大约 90K 条消息,其中可能有 1000 条(随机值)消息,RID 为 1,日期为 2020 年 1 月 1 日。
{"RID": "1" , "Date": "2020-01-01", ....}

我想在 1 分钟后仅消费 1000 条消息中的 1 条(随机 1000 条中的任意一条)。

注意:该主题有 3 个分区。

【问题讨论】:

  • 我希望我现在清楚了。谢谢指出

标签: apache-kafka apache-kafka-streams


【解决方案1】:

你想要的似乎是不可能的。代理无法根据您的业务逻辑传递消息,但它们只能传递所有条消息。

但是,您可以实现客户端缓存以相应地“去重”消息,并且在“去重”之后仅处理消息的子集。

【讨论】:

  • 明白了,谢谢,我正在探索服务器端是否有可用于优化的东西。
  • 谢谢@matthias,你能分享一下代码吗?我没有使用 confluent spring-cloud-stream 和 kafka-streams binder。
【解决方案2】:

我不完全确定您的问题,但您似乎需要compaction log

它将从主题中删除最旧的消息,只需要为主题配置压缩并使用 RID 记录作为标识符。

希望对你有帮助

【讨论】:

  • 但是我有一个持续运行的消费者,即使在压缩之前也会消耗重复的消息
猜你喜欢
  • 2019-08-05
  • 2021-05-24
  • 1970-01-01
  • 2019-11-15
  • 1970-01-01
  • 2020-03-30
  • 1970-01-01
  • 2021-06-17
  • 1970-01-01
相关资源
最近更新 更多