Flink 1.14.0编程配置S3参数连接MinIO时认证不生效问题
Flink 1.14 编程配置S3对接MinIO不生效问题修复
问题根因
你遇到的配置不生效问题由两个原因共同导致:
1. 配置键名不符合规范
Flink 内置的Hadoop S3A文件系统实现,对应的凭证提供方配置键为s3.credentials-provider,而非你使用的s3.aws.credentials.provider,错误的键名导致配置无法被识别。
2. 自定义配置未透传给Hadoop运行时
Flink 1.14 版本中,直接传入StreamExecutionEnvironment.createLocalEnvironmentWithWebUI的Flink配置,不会自动同步给底层初始化S3连接的Hadoop上下文,StreamingFileSink读取S3配置时默认只会读取Hadoop全局配置,因此拿不到你传入的自定义凭证配置。
修复方案
方案1:调整配置并显式同步Hadoop配置
修正配置键名,同时把S3相关配置同步到Hadoop全局配置中,完整代码如下:
import org.apache.hadoop.conf.Configuration as HadoopConfiguration import org.apache.flink.configuration.Configuration import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment val flinkConf = new Configuration() // Flink层面S3配置 flinkConf.setString("s3.endpoint", "http://127.0.0.1:9000") flinkConf.setString("s3.credentials-provider", "org.apache.hadoop.fs.s3a.AnonymousAWSCredentialsProvider") flinkConf.setString("s3.path.style.access", "true") // MinIO必须开启路径风格访问,否则会将bucket识别为子域名前缀 // 同步配置到Hadoop层面,供S3A文件系统读取 val hadoopConf = new HadoopConfiguration() hadoopConf.set("fs.s3a.endpoint", "http://127.0.0.1:9000") hadoopConf.set("fs.s3a.aws.credentials.provider", "org.apache.hadoop.fs.s3a.AnonymousAWSCredentialsProvider") hadoopConf.set("fs.s3a.path.style.access", "true") val env = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(flinkConf) // 绑定全局参数到运行环境 env.getConfig.setGlobalJobParameters(flinkConf)
方案2:本地调试使用配置文件
如果是本地测试场景,也可以直接把上述fs.s3a开头的配置项写入项目资源目录下的core-site.xml文件,Hadoop会自动读取该文件的配置,不需要在代码中硬编码。
额外注意事项
- 请确保项目classpath中引入的是
flink-s3-fs-hadoop依赖,不要和flink-s3-fs-presto依赖混用,两个实现的配置键规则不通用 - MinIO默认不支持虚拟主机风格访问,必须开启路径风格访问参数,否则会出现bucket解析错误
内容的提问来源于stack exchange,提问作者synapse
相关产品推荐
相关产品推荐

