Apache Flink S3文件系统凭证验证失败求助
问题:Flink运行时设置S3凭证无法通过校验
尝试通过Flink从Amazon S3读取CSV文件,在代码中运行时设置S3凭证,但始终无法通过凭证校验,相关代码及报错日志如下:
相关代码
object AwsS3CSVTest { def main(args: Array[String]): Unit = { val conf = new Configuration(); conf.setString("fs.s3a.access.key", "***") conf.setString("fs.s3a.secret.key", "***") val env = ExecutionEnvironment.createLocalEnvironment(conf) val datafile = env.readCsvFile("s3a://anybucket/anyfile.csv") .ignoreFirstLine() .fieldDelimiter(";") .types(classOf[String], classOf[String], classOf[String], classOf[String], classOf[String], classOf[String]) datafile.print() } }
报错日志
00:49:55.558|DEBUG| o.a.h.f.s.AWSCredentialProviderList No credentials from TemporaryAWSCredentialsProvider: org.apache.hadoop.fs.s3a.auth.NoAwsCredentialsException: Session credentials in Hadoop configuration: No AWS Credentials 00:49:55.558|DEBUG| o.a.h.f.s.AWSCredentialProviderList No credentials from SimpleAWSCredentialsProvider: org.apache.hadoop.fs.s3a.auth.NoAwsCredentialsException: SimpleAWSCredentialsProvider: No AWS credentials in the Hadoop configuration 00:49:55.558|DEBUG| o.a.h.f.s.AWSCredentialProviderList No credentials provided by EnvironmentVariableCredentialsProvider: 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)) 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))
解决办法
1. 修正配置参数优先级
显式指定凭证提供器,强制Flink使用代码中配置的静态凭证:
val conf = new Configuration(); conf.setString("fs.s3a.access.key", "你的访问密钥ID") conf.setString("fs.s3a.secret.key", "你的秘密访问密钥") // 强制使用SimpleAWSCredentialsProvider读取配置中的凭证 conf.setString("fs.s3a.aws.credentials.provider", "org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider") val env = ExecutionEnvironment.createLocalEnvironment(conf)
2. 校验凭证与配置正确性
- 确认
fs.s3a.access.key和fs.s3a.secret.key拼写无误,注意是s3a而非s3或其他变体; - 替换
***为真实有效的AWS密钥,确保没有多余空格、换行或特殊字符; - 测试密钥权限:用AWS CLI执行
aws s3 cp s3://anybucket/anyfile.csv ./test.csv,如果能成功下载,说明密钥权限正常。
3. 排查依赖兼容性
- 确保项目中
hadoop-aws、hadoop-common的版本与Flink依赖的Hadoop版本匹配(比如Flink 1.17默认依赖Hadoop 3.3.4); - 通过
mvn dependency:tree检查依赖树,排除冲突的Hadoop版本。
4. 备选:通过环境变量传递凭证
如果代码配置仍不生效,可直接设置系统环境变量:
- 运行程序前执行:
或在代码中添加:export AWS_ACCESS_KEY_ID="你的访问密钥ID" export AWS_SECRET_ACCESS_KEY="你的秘密访问密钥"System.setProperty("AWS_ACCESS_KEY_ID", "你的访问密钥ID") System.setProperty("AWS_SECRET_ACCESS_KEY", "你的秘密访问密钥")
内容的提问来源于stack exchange,提问作者kkurt
相关产品推荐
相关产品推荐

