【问题标题】:Read from multiple Kafka partitions and recover original order从多个Kafka分区读取并恢复原始顺序
【发布时间】:2022-02-01 21:04:55
【问题描述】:

我正在开发一个使用 Kafka 作为分布式提交日志的系统。单线程 Kafka 生产者接收来自外部的请求,处理它们并将结果写入具有 4096 个分区的主题。分区的数量是根据下游消费者的需求来选择的。生产者的内部状态会随着它接收到新的请求而变化,它会不时保存状态快照。

在极少数情况下,当生产者需要恢复时,它会读取快照,然后需要按照生成消息的顺序从 Kafka 主题中读取消息。我知道这不是 Kafka 设计的工作方式。但由于这是一种特殊且罕见的情况,我想知道我是否可以一次从每个分区读取一批,在内存中对它们进行排序,然后应用于快照以最终获得最新状态?

编辑:要记住的事情。 1.所有产生的消息都带有序列号,所以我可以订购它们。 2. Producer 在设计上是单线程的。

【问题讨论】:

  • 我不确定我是否理解“来自每个分区的批次”是什么意思。当然,您可以将分区分配给消费者(不是生产者)并轮询一次,然后转到下一个分区并重复...
  • @OneCricketeer 因为我想将 Kafka 用作“提交日志”,所以我需要 a) 写入它,b) 在恢复时从中读取。通常我的生产者会写,但当它恢复时,它变成了消费者。真正的问题是,是否有一种可靠的方法可以从多个分区中读取数据,并以某种方式按照生产者的顺序获取记录,假设每条记录都有一个序列。
  • 所有数据已经​​有一个偏移序列。但是跨多个分区排序并不是 kafka 鼓励的模式。我会指出 Kafka Streams KTables 已经可以满足您的要求

标签: apache-kafka event-sourcing


【解决方案1】:

您可以手动实现这一点,但它不会是惯用的或特别高效的(即,您应该考虑使用为事件溯源而设计的数据存储)。

基本想法是让您的生产者进程维护分区的偏移量作为快照的一部分(也可以使用消费者组,但应注意确保消费者在保留期内有效偏移量主题,否则这个生产者“太可靠”的情况可能会导致令人讨厌的意外)然后让消费者寻找这些分区偏移量,从每个分区读取 N 条消息(例如通过poll),取从合并的 4096*N 消息中连续运行最低的序列号,然后迭代,直到达到每个分区的最大偏移量。

您需要小心 Kafka 的消息重复保证:至少您有要重复数据删除的序列号。

其性能取决于您拍摄快照的频率。

【讨论】:

  • 感谢您的回答。为什么你认为它不会是高性能的?消费者尽可能快地轮询数据,唯一的事情是他们只有在获得连续的数据块时才处理数据。
  • 与线性存储消息的设置相比(可能有一个下游进程将线性消息用于恢复并将它们扇出到 4096 个分区,尽管我真的怀疑是否需要 4096 个分区:如果您在一个主题中考虑超过两位数的分区,我会认真地重新审视您如何构建消费者方面)这会有点慢。
猜你喜欢
  • 1970-01-01
  • 2016-01-17
  • 2021-07-30
  • 1970-01-01
  • 1970-01-01
  • 2015-07-01
  • 2021-03-25
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多