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

Flink 1.20.2中S3作为Source可认证但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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 07:45:54