Elasticsearch迁移期间如何暂停/恢复application模式Flink作业
可落地实现方案
方案1:无停机双写切换(优先推荐,零数据不一致风险)
不需要暂停Flink作业,从根源上规避新旧索引写入时间差的问题,是生产环境ES升级/索引重建的标准操作流程:
- 提前创建好适配新映射的ES索引,配置好分片、副本等参数
- 调整Flink作业的ES Sink逻辑,开启双写:同一条数据用相同文档ID同时写入旧索引、新索引,保证双写阶段两个索引的新增数据完全一致
- 启动存量数据迁移任务,把旧索引的历史数据同步到新索引,迁移时按文档ID做幂等判断:如果新索引里已经存在对应ID的文档(即双写阶段已经写入的新数据),直接跳过,避免旧数据覆盖新数据
- 全量迁移完成后做简单校验:对比新旧索引文档总数、抽样核对字段值,确认数据一致
- 原子切换ES读别名,把业务读流量切到新索引,观察读写正常后,调整Flink Sink配置停止写入旧索引,后续下线旧索引即可
- 整个流程没有服务停机窗口,不需要卡时间点追增量差,稳定性最高。
方案2:Flink REST API 暂停/恢复作业(适合可接受秒级停机场景)
Application模式部署的Flink自带全套REST操作接口,不需要操作K8s Pod,只要能访问Flink Web暴露的HTTP端口就能调用,完全适配权限约束场景:
暂停作业步骤
- 先查询运行中的作业列表,拿到目标同步作业的ID:
GET /jobs
返回结果里对应作业的jid字段就是作业ID。
2. 调用带savepoint的优雅停止接口:
POST /jobs/{jobId}/stop Content-Type: application/json { "drain": false, "targetDirectory": "替换为你集群配置的savepoint存储路径,比如hdfs:///flink/sp/" }
关键注意点:
drain参数必须设为false。如果设为true,Kafka源端会把消费位点直接推进到当前所有分区的最新位置,恢复作业时会漏掉暂停期间的Kafka消息;设为false时,作业会在当前消费位点触发savepoint后完全停止,位点信息完整持久化,不会丢数据。
- 等接口返回savepoint的具体存储路径,就代表作业已经暂停,此时不会有新数据写入旧ES索引,可以正常做索引映射修改、存量数据迁移操作。
恢复作业步骤
- ES侧操作全部完成后,提前把作业的ES Sink配置改成指向新索引/别名,调用作业运行接口,从之前保存的savepoint启动即可:
POST /jars/{已上传的作业jar包ID}/run Content-Type: application/json { "savepointPath": "上一步返回的savepoint全路径", "allowNonRestoredState": false, "entryClass": "你的作业主类全限定名", "programArgs": "替换为修改后的作业启动参数,包含新的ES索引/别名配置" }
- 作业启动后会自动从暂停时的Kafka位点开始消费,把暂停期间积压的所有消息写入新ES索引,不会出现数据丢失,也不会有旧索引额外写入导致的不一致问题。
方案3:运行时热控速实现逻辑暂停(适合短时间操作场景)
如果不想停作业,可以直接调用Flink的动态配置接口调整Kafka源端的消费限速,把速率设为0即可实现逻辑暂停,不需要重启作业:
- 先查询作业的算子列表,找到Kafka Source对应的算子ID:
GET /jobs/{jobId}/vertices
- 调用算子配置更新接口,把源端消费速率设为0:
PATCH /jobs/{jobId}/vertices/{source算子ID}/config Content-Type: application/json { "ratelimit.records-per-second": "0" }
配置生效后Kafka源端会停止拉取新消息,不会有新数据写入ES,此时可以做迁移操作。等操作完成后,再把这个参数改回正常的消费限速值,作业就会自动恢复消费。
注意:这个方案仅适合几小时内的短时间暂停,如果暂停时长超过Kafka集群配置的消息保留时间,会因为历史消息被清理导致数据丢失,长时间操作优先选带savepoint的停止/恢复方案。
内容的提问来源于stack exchange,提问作者user4695271
相关产品推荐
相关产品推荐

