【发布时间】:2019-06-03 04:28:14
【问题描述】:
我们使用 kafka 作为管道的来源。我想将现有状态从生产环境转移到新环境。我的问题是新环境中的偏移量会发生什么?由于我们从生产中获取了保存点,并且偏移量保存在保存点中,这是否意味着在新环境中,作业将开始使用来自生产的偏移量的消息,或者它实际上会从一个新的消息开始,比如作为新的消费者?
【问题讨论】:
标签: apache-kafka apache-flink flink-streaming
我们使用 kafka 作为管道的来源。我想将现有状态从生产环境转移到新环境。我的问题是新环境中的偏移量会发生什么?由于我们从生产中获取了保存点,并且偏移量保存在保存点中,这是否意味着在新环境中,作业将开始使用来自生产的偏移量的消息,或者它实际上会从一个新的消息开始,比如作为新的消费者?
【问题讨论】:
标签: apache-kafka apache-flink flink-streaming
新作业中的偏移量将从保存点中存储的偏移量开始,前提是您从保存点重新启动新作业,如下所示:
$ bin/flink run -s :savepointPath [:runArgs]
相关文档包括本节关于Kafka Consumers Start Position Configuration 的最后一段,其中指出
请注意,当作业从故障中自动恢复或使用保存点手动恢复时,这些开始位置配置方法不会影响开始位置。在还原时,每个 Kafka 分区的起始位置由保存点或检查点中存储的偏移量确定...
以及关于Resuming from Savepoints 的部分。
【讨论】: