Flink 1.12.2 本地修改S3来源Savepoint因元数据绝对路径报错如何解决
问题根因
你拷贝到本地的Savepoint的_metadata文件中存储的所有状态文件路径都是S3上的绝对路径,Flink State Processor API加载Savepoint时,会直接按_metadata中记录的路径读取文件,而你本地没有对应路径的文件,所以抛出FileNotFoundException。
解决方案
方案1:直接读写S3 Savepoint(推荐)
不需要把Savepoint拷到本地,直接通过Flink的S3文件系统插件读写S3上的状态,完全规避路径匹配问题:
- 项目引入S3文件系统依赖:
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-s3-fs-hadoop</artifactId> <version>1.12.2</version> </dependency>
- 给执行环境配置S3访问凭证:
Configuration conf = new Configuration(); conf.setString("s3.access-key", "你的S3访问AK"); conf.setString("s3.secret-key", "你的S3访问SK"); // 自建S3兼容存储需额外配置endpoint conf.setString("s3.endpoint", "你的S3服务地址"); ExecutionEnvironment executionEnvironment = ExecutionEnvironment.createLocalEnvironment(conf);
- 直接加载S3路径的原Savepoint,修改后可以选择写本地或者直接写S3:
BootstrapTransformation<AccountRegistrationInformation> transformation = OperatorTransformation .bootstrapWith(accountDataSet) .keyBy(acc -> acc.getBrand() + "-" + acc.getAccountId()) .transform(new AccountRegistrationBootstrapper()); // 直接读取S3上的原Savepoint Savepoint.load(executionEnvironment, "s3://<redacted>/savepoint-c680a3-c178150a8b8d", new MemoryStateBackend()) .removeOperator("registration-processor") .withOperator("registration-processor", transformation) // 填S3路径可直接把修改后的Savepoint写入S3,无需手动上传 .write("s3://<redacted>/transformed-savepoint"); executionEnvironment.execute();
方案2:基于本地拷贝的Savepoint修改
如果网络限制无法直接访问S3,可按以下方式处理:
Flink 1.12版本的_metadata是二进制序列化格式,无法直接文本修改路径,你可以直接在本地创建和报错路径完全匹配的目录层级,将所有状态数据文件放到对应目录下即可。比如报错提示找不到<redacted>\savepoint-c680a3-c178150a8b8d\32c44059-xxx,就在你运行代码的盘符下创建同名的层级目录,把所有数据文件放入对应目录,Flink就能正常读取。
修改完成后新生成的Savepoint的_metadata里存的是本地路径,上传到S3之前,需要重新读取本地生成的Savepoint,用方案1的方式直接写入S3路径,自动生成带S3路径的_metadata,无需手动修改二进制文件。
注意事项
- 所有修改操作都基于原Savepoint的副本执行,不要改动原Savepoint文件,避免原状态损坏。
- 替换状态的算子
registration-processor必须和原作业中算子设置的uid完全一致,否则会出现状态匹配失败的问题。 - 修改完成后建议先执行
flink run -s <新Savepoint路径> --dry-run验证Savepoint可用性,避免线上作业启动失败。
内容的提问来源于stack exchange,提问作者Mihai Banu
相关产品推荐
相关产品推荐

