You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.18 17:05:44