Flink 1.20.2中S3作为Source可认证但Sink认证失败求助
问题:Flink 1.20.2 写入S3 Sink时认证失败
我用Flink 1.20.2编写批处理作业,从S3的Parquet文件读取多租户数据,先提取租户ID列表,再按租户拆分写入单独的S3路径。读取操作正常,但写入时报权限拒绝,认证失败。
作业代码
public class FlinkParquetJob { private static void runJob() throws Exception { System.out.println(Instant.now().toString()); EnvironmentSettings settings = EnvironmentSettings.newInstance().inBatchMode().build(); TableEnvironment tEnv = TableEnvironment.create(settings); var schema = Helper.getSchema(); tEnv.createTemporaryTable( "source", TableDescriptor.forConnector("filesystem") .schema(schema) .option("path", "s3a://data/input/input.parquet") .option("format", "parquet") .build() ); List<Long> tenants = new ArrayList<>(); try (var rows = tEnv.executeSql("select distinct tenant_id from source").collect()) { while (rows.hasNext()) { Row row = rows.next(); Long tenant = (Long) row.getField("tenant_id"); tenants.add(tenant); } } System.out.println("Found " + tenants.size() + " tenants. Creating statementset..."); for (Long tenant : tenants) { String sinkTableName = String.format("tenant_%s", tenant.toString()); tEnv.createTemporaryTable(sinkTableName, TableDescriptor.forConnector("filesystem") .schema(schema) .option("path", String.format("s3a://data/output/tenant_%s", tenant.toString())) .option("format", "json") .build() ); tEnv.sqlQuery(String.format("select * from source where tenant_id = %s", tenant.toString())) .executeInsert(sinkTableName); } } public static void main(String[] args) throws Exception { runJob(); } }
错误信息
Caused by: java.nio.file.AccessDeniedException: s3a://data/output/tenant_24236234: org.apache.hadoop.fs.s3a.auth.NoAuthWithAWSException: No AWS Credentials provided by DynamicTemporaryAWSCredentialsProvider TemporaryAWSCredentialsProvider SimpleAWSCredentialsProvider EnvironmentVariableCredentialsProvider IAMInstanceCredentialsProvider : com.amazonaws.SdkClientException: Unable to load AWS credentials from environment variables (AWS_ACCESS_KEY_ID (or AWS_ACCESS_KEY) and AWS_SECRET_KEY (or AWS_SECRET_ACCESS_KEY)) at org.apache.hadoop.fs.s3a.S3AUtils.translateException(S3AUtils.java:212) at org.apache.hadoop.fs.s3a.S3AUtils.translateException(S3AUtils.java:175) at org.apache.hadoop.fs.s3a.S3AFileSystem.s3GetFileStatus(S3AFileSystem.java:3799) at org.apache.hadoop.fs.s3a.S3AFileSystem.innerGetFileStatus(S3AFileSystem.java:3688) at org.apache.hadoop.fs.s3a.S3AFileSystem.lambda$exists$34(S3AFileSystem.java:4703) at org.apache.hadoop.fs.statistics.impl.IOStatisticsBinding.lambda$trackDurationOfOperation$5(IOStatisticsBinding.java:499) at org.apache.hadoop.fs.statistics.impl.IOStatisticsBinding.trackDuration(IOStatisticsBinding.java:444) at org.apache.hadoop.fs.s3a.S3AFileSystem.trackDurationAndSpan(S3AFileSystem.java:2337) at org.apache.hadoop.fs.s3a.S3AFileSystem.trackDurationAndSpan(S3AFileSystem.java:2356) at org.apache.hadoop.fs.s3a.S3AFileSystem.exists(S3AFileSystem.java:4701) at org.apache.flink.fs.s3hadoop.common.HadoopFileSystem.exists(HadoopFileSystem.java:165) at org.apache.flink.core.fs.PluginFileSystemFactory$ClassLoaderFixingFileSystem.exists(PluginFileSystemFactory.java:148) at org.apache.flink.connector.file.table.FileSystemOutputFormat.createStagingDirectory(FileSystemOutputFormat.java:110) ... 46 more Caused by: org.apache.hadoop.fs.s3a.auth.NoAuthWithAWSException: No AWS Credentials provided by DynamicTemporaryAWSCredentialsProvider TemporaryAWSCredentialsProvider SimpleAWSCredentialsProvider EnvironmentVariableCredentialsProvider IAMInstanceCredentialsProvider : com.amazonaws.SdkClientException: Unable to load AWS credentials from environment variables (AWS_ACCESS_KEY_ID (or AWS_ACCESS_KEY) and AWS_SECRET_KEY (or AWS_SECRET_ACCESS_KEY)) at org.apache.hadoop.fs.s3a.AWSCredentialProviderList.getCredentials(AWSCredentialProviderList.java:216) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.getCredentialsFromContext(AmazonHttpClient.java:1295) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.runBeforeRequestHandlers(AmazonHttpClient.java:869) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.doExecute(AmazonHttpClient.java:818) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.executeWithTimer(AmazonHttpClient.java:805) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.execute(AmazonHttpClient.java:779) at com.amazonaws.http.AmazonHttpClient$RequestExecutor.access$500(AmazonHttpClient.java:735) at com.amazonaws.http.AmazonHttpClient$RequestExecutionBuilderImpl.execute(AmazonHttpClient.java:717) at com.amazonaws.http.AmazonHttpClient.execute(AmazonHttpClient.java:581) at com.amazonaws.http.AmazonHttpClient.execute(AmazonHttpClient.java:559) at com.amazonaws.services.s3.AmazonS3Client.invoke(AmazonS3Client.java:5593) at com.amazonaws.services.s3.AmazonS3Client.getBucketRegionViaHeadRequest(AmazonS3Client.java:6574) at com.amazonaws.services.s3.AmazonS3Client.fetchRegionFromCache(AmazonS3Client.java:6546) at com.amazonaws.services.s3.AmazonS3Client.invoke(AmazonS3Client.java:5578) at com.amazonaws.services.s3.AmazonS3Client.invoke(AmazonS3Client.java:5540) at com.amazonaws.services.s3.AmazonS3Client.getObjectMetadata(AmazonS3Client.java:1422) at org.apache.hadoop.fs.s3a.S3AFileSystem.lambda$getObjectMetadata$10(S3AFileSystem.java:2545) at org.apache.hadoop.fs.s3a.Invoker.retryUntranslated(Invoker.java:414) at org.apache.hadoop.fs.s3a.Invoker.retryUntranslated(Invoker.java:377) at org.apache.hadoop.fs.s3a.S3AFileSystem.getObjectMetadata(S3AFileSystem.java:2533) at org.apache.hadoop.fs.s3a.S3AFileSystem.getObjectMetadata(S3AFileSystem.java:2513) at org.apache.hadoop.fs.s3a.S3AFileSystem.s3GetFileStatus(S3AFileSystem.java:3776) ... 56 more Caused by: com.amazonaws.SdkClientException: Unable to load AWS credentials from environment variables (AWS_ACCESS_KEY_ID (or AWS_ACCESS_KEY) and AWS_SECRET_KEY (or AWS_SECRET_ACCESS_KEY)) at com.amazonaws.auth.EnvironmentVariableCredentialsProvider.getCredentials(EnvironmentVariableCredentialsProvider.java:66) at org.apache.hadoop.fs.s3a.AWSCredentialProviderList.getCredentials(AWSCredentialProviderList.java:177) ... 77 more
项目依赖
implementation "org.apache.flink:flink-core:1.20.2" implementation "org.apache.flink:flink-python:1.20.2" implementation "org.apache.flink:flink-java:1.20.2" implementation "org.apache.flink:flink-clients:1.20.2" implementation "org.apache.flink:flink-streaming-java_2.12:1.9.3" implementation "org.apache.flink:flink-table-api-java-bridge_2.12:1.9.3" implementation "org.apache.flink:flink-table-planner-loader:1.20.2" runtimeOnly "org.apache.flink:flink-table-runtime:1.20.2" implementation "org.apache.flink:flink-s3-fs-hadoop:1.20.2" implementation "org.apache.iceberg:iceberg-flink-runtime-1.20:1.9.2" implementation platform ("software.amazon.awssdk:bom:2.20.135") implementation "software.amazon.awssdk:s3" implementation "software.amazon.awssdk:s3-transfer-manager"
已尝试的解决步骤
- 在Flink容器中设置了
AWS_SECRET_ID和AWS_SECRET_KEY环境变量 - 在JobManager和TaskManager的
flink-conf.yaml中添加配置:fs.s3a.access.key: *** fs.s3a.secret.key: *** fs.s3a.aws.credentials.provider: org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider - 启动作业时通过命令行指定参数:
flink run -m 127.0.0.1:10111 -c com.company.FlinkParquetJob program.jar -Dfs.s3a.access.key=*** -Dfs.s3a.secret.key=***
问题排查与解决方案
1. 修复依赖版本冲突
你的依赖中混合了Flink 1.20.2和1.9.3的组件,这会导致类加载冲突,尤其是S3认证逻辑。必须统一所有Flink依赖版本为1.20.2:
implementation "org.apache.flink:flink-core:1.20.2" implementation "org.apache.flink:flink-python:1.20.2" implementation "org.apache.flink:flink-java:1.20.2" implementation "org.apache.flink:flink-clients:1.20.2" implementation "org.apache.flink:flink-streaming-java_2.12:1.20.2" implementation "org.apache.flink:flink-table-api-java-bridge_2.12:1.20.2" implementation "org.apache.flink:flink-table-planner-loader:1.20.2" runtimeOnly "org.apache.flink:flink-table-runtime:1.20.2" implementation "org.apache.flink:flink-s3-fs-hadoop:1.20.2" implementation "org.apache.iceberg:iceberg-flink-runtime-1.20:1.9.2" implementation platform ("software.amazon.awssdk:bom:2.20.135") implementation "software.amazon.awssdk:s3" implementation "software.amazon.awssdk:s3-transfer-manager"
2. 修正环境变量名称
错误日志显示,Hadoop的EnvironmentVariableCredentialsProvider只识别AWS_ACCESS_KEY_ID和AWS_SECRET_ACCESS_KEY,你设置的AWS_SECRET_ID和AWS_SECRET_KEY名称不匹配,需修改:
export AWS_ACCESS_KEY_ID=你的访问密钥ID export AWS_SECRET_ACCESS_KEY=你的私有访问密钥
3. 优化作业逻辑,避免客户端-集群凭证传递问题
当前代码在客户端(main方法)提前收集租户列表,这部分能读取到环境变量,但后续executeInsert提交到集群执行时,TaskManager可能无法获取客户端的环境变量。改用Flink动态分区写入,无需手动遍历租户:
public class FlinkParquetJob { private static void runJob() throws Exception { System.out.println(Instant.now().toString()); EnvironmentSettings settings = EnvironmentSettings.newInstance().inBatchMode().build(); TableEnvironment tEnv = TableEnvironment.create(settings); var schema = Helper.getSchema(); tEnv.createTemporaryTable( "source", TableDescriptor.forConnector("filesystem") .schema(schema) .option("path", "s3a://data/input/input.parquet") .option("format", "parquet") .build() ); // 按tenant_id分区自动拆分文件 tEnv.createTemporaryTable( "sink", TableDescriptor.forConnector("filesystem") .schema(schema) .option("path", "s3a://data/output/") .option("format", "json") .option("partition.column", "tenant_id") // 控制每个分区生成单个文件 .option("sink.rolling-policy.file-size", "1024MB") .option("sink.rolling-policy.rollover-interval", "1h") .build() ); tEnv.executeSql("insert into sink select * from source"); } public static void main(String[] args) throws Exception { runJob(); } }
4. 确保集群配置正确生效
检查flink-conf.yaml配置格式,必须分行编写,否则会解析失败:
fs.s3a.access.key: 你的访问密钥ID fs.s3a.secret.key: 你的私有访问密钥 fs.s3a.aws.credentials.provider: org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider
修改后重启JobManager和TaskManager,确保配置生效。
内容的提问来源于stack exchange,提问作者LuckyGambler
相关产品推荐
相关产品推荐

