【问题标题】:Kafka producer - How to change a topic without down-time and preserving message ordering?Kafka 生产者 - 如何在不停机和保留消息顺序的情况下更改主题?
【发布时间】: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

正式任务

v1v2 成为主题。目前,我将format_v1 格式的消息生成到主题v1 中。我想将format_v2 格式的消息生成到主题v2 中。切换应该在我可以选择的某个时间发生。

换句话说,在某个时刻,生产者的所有实例都停止向v1发送消息,并开始向v2发送消息;因此v1 中的最后一条消息m1v2m2 的第一条消息之前生成。

详情

我有一个想法,我可以为主题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_v1format_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。

我想在数据库中手动操作(通过一些休息端点等)来更改标志;生产者在生产消息之前总是检查标志。当标志显示v2true 时,生产者将开始将消息写入v2。但是,如果在某个时刻标志为假,一个生产者开始向v1 生成消息,然后标志已更改并且另一个生产者在前一个生产者完成对v1 的生产之前已将消息发送到v2

【问题讨论】:

    标签: avro messaging kafka-producer-api


    【解决方案1】:

    您可以接受只有一名制作人在场吗?

    在这种情况下,您可以将您的想法与标志一起使用:

    1. 关闭所有生产者p2,p3,...,pn 除了p1
    2. p1 单独写信给v1
    3. 将标志切换为v2,因此p1 结束对v1 的最后一次写入并开始写入v2
    4. 现在没有人写信给v1
    5. 开始你的其他生产者p2,p3,...,pn
    6. 现在每个生产者都因为活动标志写信给v2,但仍然没有人写信给v1

    【讨论】:

    • 我同意,这是可能的,谢谢。我会做更多的研究,看看也许还有其他建议。
    • 虽然有一个缺点。看,我不是要批评你的回答,我只是希望每个阅读的人都能学到。上述解决方案的优点是可行且不复杂(可以完成并且很清楚如何!)。另一方面,假设您有多个环境:开发、暂存、测试、预生产、生产。这意味着有工作需要完成:for each environment:stop producers except oneswitch feature flagstart producers again。试想一下,如何自动化所有这些工作(或减少手工工作量)。在这种情况下,这将是惊人的!
    • 是的,你是对的。我不建议手动执行这些步骤。将您计划多次执行的所有操作自动化,以节省时间并防止错误。在我天真的想象中,这样的自动化脚本很容易实现,因为每个步骤都很简单。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-06-12
    • 2020-12-11
    • 1970-01-01
    • 2021-05-28
    • 1970-01-01
    • 1970-01-01
    • 2017-04-29
    相关资源
    最近更新 更多