【问题标题】:Kafka - consuming messages based on timestampKafka - 基于时间戳消费消息
【发布时间】:2020-02-18 11:32:47
【问题描述】:

我是 Kafka 的新手,但需要实现消费者基于时间戳从特定主题消费的逻辑。另一个用例也是让我能够在特定时间范围内消费(例如从 10:00 到 10:20)。该范围将始终可除以 5 分钟 - 这意味着我不需要从例如 10:00 到 10:04 消费)。我想的逻辑如下:

  1. 创建一个存储时间戳和 Kafka messageId (timestamp | id) 的表
  2. 创建一个每 5 分钟执行以下操作的控制台\服务:
  3. 获取某个主题的所有分区
  4. 查询所有分区的最小偏移值(一个起点)
  5. 将偏移量和时间戳存储在表中

获取一个主题的所有分区

现在,如果一切正常,我应该在桌子上放一些类似的东西:

10:00     | 0
10:05     | 100
10:10     | 200
HH: mm     | (some number)

现在有了这个,我可以随时启动消费者并知道我应该能够消费我需要的偏移量。 它看起来是正确的还是我在某个地方出现了缺陷?或者也许有更好的方法来实现所需的结果?任何想法或建议将不胜感激。

P.S.:我的一位同事建议使用分区并分别处理每个分区...这意味着如果我有一个主题并且副本数例如为 5 - 那么我需要为我的主题保存 5 次偏移量对于每个间隔(每个分区一次)。然后消费者还需要考虑分区并根据我为每个分区获得的偏移量进行消费。但这会包含我试图避免的额外复杂性......

提前致谢!

BR, 迈克

【问题讨论】:

    标签: apache-kafka kafka-consumer-api


    【解决方案1】:

    不需要桌子。

    您可以使用 Consumer 实例的 seek 方法将所有分区移动到该分区定义的偏移量。

    分区可能会起作用... 5 分钟消息间隔的 12 个分区

    我认为复制不能解决您的问题。

    【讨论】:

      猜你喜欢
      • 2020-07-08
      • 1970-01-01
      • 1970-01-01
      • 2019-05-02
      • 1970-01-01
      • 1970-01-01
      • 2018-03-06
      • 2019-07-06
      相关资源
      最近更新 更多