【发布时间】:2020-06-12 20:30:38
【问题描述】:
这个问题是关于架构和 kafka 主题迁移的。
原始问题:没有向后兼容性的模式演变。
https://docs.confluent.io/current/schema-registry/avro.html
我正在请求社区给我一个建议或分享文章,我可以从中获得灵感,也许可以想出解决我的问题的方法。也许有一种架构或流模式。没有必要给我一个特定语言的解决方案;只要给我一个我可以去的方向......我的问题很大,对于以后想要的人来说可能会很有趣
- a) 更改消息格式并将消息生成新主题。
- b) 停止向一个主题生成消息,并“立即”开始向另一个主题生成消息;换句话说,一旦产生了
v2中的消息,就不会将新消息附加到v1。
问题
我正在更改消息格式,与以前的版本不兼容。为了不破坏现有消费者,我决定为新主题生成消息。
向上施法者的想法
我读过关于上施法者的文章。
https://docs.axoniq.io/reference-guide/operations-guide/production-considerations/versioning-events
正式任务
让v1 和v2 成为主题。目前,我将format_v1 格式的消息生成到主题v1 中。我想将format_v2 格式的消息生成到主题v2 中。切换应该在我可以选择的某个时间发生。
换句话说,在某个时刻,生产者的所有实例都停止向v1发送消息,并开始向v2发送消息;因此v1 中的最后一条消息m1 在v2 中m2 的第一条消息之前生成。
详情
我有一个想法,我可以为主题v1 生成消息,有一个订阅v1 的kafka steam up-caster 并将转换后的消息推送到v2。假设转换器(当然是我的情况)能够将format_v1 的消息转换为format_v2 而不会出错。
如上面关于 avro 模式演变的链接中所述,当我添加一个上施法者并将消息生成到 v1 时,我已经将 v1 的所有消费者更改为 v2。
现在,一个棘手的部分。我们有两个要求:
1.没有生产停机时间。
2. 保持消息顺序。
意思是:
1) 我们不允许丢失消息;客户可以随时使用我们的系统,因此我们的系统应随时生成消息。
2) 我们正在运行生产者的多个实例。在某个时刻,可能(可能)有生产者可以将格式为 format_v1 的消息生成到主题 v1 中,并且某些实例会生成格式为 format_v2 的消息到主题 v2 中。
众所周知,kafka 不保证不同分区和主题的消息排序。
我可以通过使用与 v1 相同的分区选择器将消息写入 v2 来解决分区问题。或者现在,我可以想象我们只为v1 使用一个分区,为v2 使用一个分区。
我的简化和尝试
1) 我想,当我想改变生产者以将消息生成到一个新主题时,我有一个能够将消息从 v1 转换为 v2 的上施法者(kafka 流组件)没有错误。这个 kafka 流组件是可扩展的。
2) 我所有的消费者都已经切换到v2 主题。他们不断收到来自v2 的消息。此时此刻,我的生产者实例正在向主题 v1 生成消息,而上施法者的工作做得很好。
3) 为简化问题,假设现在format_v1 和format_v2 无关紧要,它们是相同的。
4) 假设我们有一个用于v1 的分区和一个用于v2 的分区。
现在我的问题是,如何从给定时间点立即切换所有生产者;所有实例都在主题 v2 中生成消息。
我的同事和卡夫卡专家告诉我,可以通过停机时间来完成
如果您依赖分区中消息的顺序,则无法在不停机的情况下切换到新版本。为了尽量减少停机时间,我们可以执行以下操作。
Upcaster 组件必须将数据写入相同的分区,并且应该尝试做出相同的偏移量。然而,这并不总是可能的,因为偏移量可能有间隙,因此必须保留旧偏移量和新偏移量之间的映射。没有所有记录,只有每个分区的最后一个批量。如果 upcaster 崩溃了,重新开始,producer 仍然不参与 v2。
启动 v2 消费者。如果和v1同一个consumer group,什么都不用做,如果有新的consumer group,根据新的offset更新Kafka中的offset。
现在 Producers 写入 v1,upcaster 转换数据,consumer 从 v2 消费
时间到了。当 upcaster 的 lag 接近 0 时,关闭 v1 producer,等到 upcaster 转换剩余的记录后,关闭 upcaster,启动 v2 producer,写入 v2 topic。
我想在数据库中手动操作(通过一些休息端点等)来更改标志;生产者在生产消息之前总是检查标志。当标志显示v2 或true 时,生产者将开始将消息写入v2。但是,如果在某个时刻标志为假,一个生产者开始向v1 生成消息,然后标志已更改并且另一个生产者在前一个生产者完成对v1 的生产之前已将消息发送到v2。
【问题讨论】:
标签: avro messaging kafka-producer-api