Flink StreamingFileSink写入S3遇AWS令牌过期问题咨询
问题描述
尝试通过s3a协议将Kafka中的数据流式写入Amazon S3,数据管道运行正常,但1小时后(与AWS令牌过期时间一致),StreamingFileSink抛出如下异常:
Caused by: com.amazonaws.services.s3.model.AmazonS3Exception: The provided token has expired. (Service: Amazon S3; Status Code: 400; Error Code: ExpiredToken; Request ID: 7YFGVQ92YT51DP0K; S3 Extended Request ID: sx6UJJ548o0wpwJbkoWJ16jKRVih3ZV9XQdbThNhq5kUU7A7yCx58tcCGELVs5tqGWaMMPfZxZM=; Proxy: webproxy) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.handleErrorResponse(AmazonHttpClient.java:1819) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.handleServiceErrorResponse(AmazonHttpClient.java:1403) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.executeOneRequest(AmazonHttpClient.java:1372) ...
使用了自定义的AWSCredentialsProvider实现,重写getCredentials方法,每15分钟从AWS获取新密钥刷新令牌。相关初始化代码如下:
StreamExecutionEnvironment env = getStreamExecutionEnvironment(); StreamingFileSink<FELEvent> sink = StreamingFileSink .forBulkFormat(<my Path Settings with basePath s3a://bucket/path/to/dir>) .withRollingPolicy(OnCheckpointRollingPolicy.build()) .withNewBucketAssigner(<My custom bucket assigner>) .build(); env.fromSource(<Kafka source>) .map(<Some operation>) .filter(<Some filtering>) .addSink(sink) .name("name").uid("uid"); env.execute("TAG");
补充:查看hadoop-fs插件代码后发现,它仅在FileSink初始化时使用提供的令牌创建一次S3对象,现寻求重新初始化的方法。
核心问题:StreamingFileSink是否会为已初始化的实例自动刷新令牌?如果不会,处理该场景的最佳方式是什么?(使用Flink 14.3版本,因ZooKeeper兼容性问题无法升级)
解决方案
1. 令牌刷新机制说明
Flink 1.14.3的StreamingFileSink本身不会自动刷新已初始化实例的AWS令牌。hadoop-fs插件在初始化阶段一次性创建S3客户端实例,后续不会主动重新获取凭证,这就是令牌过期后报错的直接原因。
2. 最佳处理方式
方式一:配置Hadoop自动刷新凭证
修改Hadoop的core-site.xml配置,让S3A客户端自动调用你的自定义凭证提供者刷新令牌:
<property> <name>fs.s3a.aws.credentials.provider</name> <value>com.your.custom.CredentialsProvider</value> </property> <property> <name>fs.s3a.session.token.refresh.period</name> <value>900s</value> <!-- 15分钟,与你的刷新周期匹配 --> </property>
这种方式依赖Hadoop S3A客户端的内置逻辑,会定期调用凭证提供者的getCredentials方法获取新凭证,无需修改Flink业务代码。
方式二:在自定义BucketAssigner中嵌入凭证刷新逻辑
既然已经使用了自定义BucketAssigner,可以在其中加入凭证刷新触发逻辑:
- 利用
getBucketId方法的调用时机,每次触发时确保自定义AWSCredentialsProvider执行刷新(前提是你的提供者已经实现了自动刷新逻辑) - 或者在BucketAssigner中持有S3客户端引用,定期重新创建客户端(需保证线程安全,适配Flink多线程运行环境)
方式三:结合Flink重启策略实现凭证重新初始化
配置Flink固定延迟重启策略,当令牌过期报错时,Flink自动重启任务,此时StreamingFileSink会重新初始化并获取新的有效令牌。这种方式实现简单,但会有短暂任务中断,适合对中断容忍度较高的场景:
env.setRestartStrategy(RestartStrategies.fixedDelayRestart( 3, // 最大重启次数 Time.seconds(10) // 重启间隔时间 ));
3. 注意事项
- 自定义
AWSCredentialsProvider必须保证线程安全,避免多线程调用时出现凭证不一致问题 - 若修改Hadoop配置,需确保Flink集群所有节点都能加载到更新后的
core-site.xml文件
内容的提问来源于stack exchange,提问作者Dimitrios

