如何通过Hashicorp Vault为K8s部署的Flink提供S3凭证?
运行时为Flink Kubernetes Operator部署的应用提供Minio凭证问题
问题描述
使用官方Flink Kubernetes Operator部署Flink流处理应用,以Minio作为状态后端。原本直接在flinkConfiguration中配置presto.s3.access-key和presto.s3.secret-key时运行正常,现在尝试通过Hashicorp Vault注入凭证,遇到以下问题:
- 注释掉静态配置项,通过Vault Agent注入凭证文件
/vault/secrets/appsecrets.yaml,并在代码中读取该文件加载到Flink配置后,出现凭证加载失败的错误:
java.io.IOException: com.amazonaws.SdkClientException: Unable to load AWS credentials from any provider in the chain: [EnvironmentVariableCredentialsProvider: 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)), SystemPropertiesCredentialsProvider: Unable to load AWS credentials from Java system properties (aws.accessKeyId and aws.secretKey), WebIdentityTokenCredentialsProvider: You must specify a value for roleArn and roleSessionName, com.amazonaws.auth.profile.ProfileCredentialsProvider@5331f738: profile file cannot be null, com.amazonaws.auth.EC2ContainerCredentialsProviderWrapper@bc0353f: Failed to connect to service endpoint: ] at com.facebook.presto.hive.s3.PrestoS3FileSystem$PrestoS3OutputStream.uploadObject(PrestoS3FileSystem.java:1278) ~[flink-s3-fs-presto-1.14.2.jar:1.14.2] at com.facebook.presto.hive.s3.PrestoS3FileSystem$PrestoS3OutputStream.close(PrestoS3FileSystem.java:1226) ~[flink-s3-fs-presto-1.14.2.jar:1.14.2] at org.apache.hadoop.fs.FSDataOutputStream$PositionCache.close(FSDataOutputStream.java:72) ~[flink-s3-fs-presto-1.14.2.jar:1.14.2] at org.apache.hadoop.fs.FSDataOutputStream.close(FSDataOutputStream.java:101) ~[flink-s3-fs-presto-1.14.2.jar:1.14.2] at org.apache.flink.fs.s3presto.common.HadoopDataOutputStream.close(HadoopDataOutputStream.java:52) ~[flink-s3-fs-presto-1.14.2.jar:1.14.2] at org.apache.flink.runtime.blob.FileSystemBlobStore.put(FileSystemBlobStore.java:80) ~[flink-dist_2.12-1.14.2.jar:1.14.2] at org.apache.flink.runtime.blob.FileSystemBlobStore.put(FileSystemBlobStore.java:72) ~[flink-dist_2.12-1.14.2.jar:1.14.2] at org.apache.flink.runtime.blob.BlobUtils.moveTempFileToStore(BlobUtils.java:385) ~[flink-dist_2.12-1.14.2.jar:1.14.2] at org.apache.flink.runtime.blob.BlobServer.moveTempFileToStore(BlobServer.java:680) ~[flink-dist_2.12-1.14.2.jar:1.14.2] at org.apache.flink.runtime.blob.BlobServerConnection.put(BlobServerConnection.java:350) [flink-dist_2.12-1.14.2.jar:1.14.2] at org.apache.flink.runtime.blob.BlobServerConnection.run(BlobServerConnection.java:110) [flink-dist_2.12-1.14.2.jar:1.14.2]
- 尝试在
docker-entrypoint.sh中将凭证追加到flink-conf.yaml,但该文件由Flink Operator从ConfigMap挂载,文件系统为只读,无法修改。
解决方案
方案1:Vault Agent直接注入环境变量
利用Flink PrestoS3FileSystem自动读取AWS标准环境变量的特性,让Vault Agent将Minio凭证注入为环境变量,无需修改代码或配置文件。
修改FlinkDeployment中JobManager和TaskManager的Pod注解(两者都需要访问Minio):
apiVersion: flink.apache.org/v1beta1 kind: FlinkDeployment metadata: name: flink-app namespace: default spec: # 其他原有配置... jobManager: podTemplate: metadata: annotations: vault.hashicorp.com/namespace: "/example/dev" vault.hashicorp.com/agent-inject: "true" vault.hashicorp.com/agent-init-first: "true" vault.hashicorp.com/role: "example-serviceaccount" vault.hashicorp.com/auth-path: auth/example # 注入Access Key到环境变量 vault.hashicorp.com/agent-inject-env-AWS_ACCESS_KEY_ID: "example/Minio#accessKey" # 注入Secret Key到环境变量 vault.hashicorp.com/agent-inject-env-AWS_SECRET_ACCESS_KEY: "example/Minio#secretKey" taskManager: podTemplate: metadata: annotations: vault.hashicorp.com/namespace: "/example/dev" vault.hashicorp.com/agent-inject: "true" vault.hashicorp.com/agent-init-first: "true" vault.hashicorp.com/role: "example-serviceaccount" vault.hashicorp.com/auth-path: auth/example vault.hashicorp.com/agent-inject-env-AWS_ACCESS_KEY_ID: "example/Minio#accessKey" vault.hashicorp.com/agent-inject-env-AWS_SECRET_ACCESS_KEY: "example/Minio#secretKey"
此方案下,Flink的S3客户端会自动从环境变量获取凭证,覆盖所有需要凭证的组件(包括BlobStore)。
方案2:自定义可写配置目录
通过挂载临时目录作为额外配置目录,让Flink加载该目录下的凭证配置,绕过只读的默认flink-conf.yaml。
- 在PodTemplate中添加emptyDir卷:
podTemplate: spec: volumes: - name: custom-conf emptyDir: {} containers: - name: flink-main-container volumeMounts: - name: custom-conf mountPath: /opt/flink/conf-custom
- 修改
docker-entrypoint.sh,将Vault凭证文件复制到自定义配置目录,并调整Flink启动参数加载多配置目录:
if [ -f '/vault/secrets/appsecrets.yaml' ]; then # 复制凭证到自定义配置目录 cp '/vault/secrets/appsecrets.yaml' /opt/flink/conf-custom/flink-conf.yaml fi # 添加-c参数,让Flink同时加载默认目录和自定义目录的配置 exec "$@" -c /opt/flink/conf:/opt/flink/conf-custom
Flink会自动合并多个配置目录的内容,自定义目录中的配置优先级更高。
方案3:代码中手动配置FileSystem凭证
如果必须通过代码加载凭证,需要确保在FileSystem初始化前就设置好凭证(因为BlobStore在JobManager启动时就会初始化,早于用户代码执行):
import org.apache.flink.configuration.Configuration import org.apache.flink.fs.s3presto.PrestoS3FileSystemFactory import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment import org.yaml.snakeyaml.Yaml import java.io.FileReader // 读取Vault中的凭证文件 def getSecretsFromFile(path: String): Map[String, String] = { new Yaml().load(new FileReader(path)).asInstanceOf[java.util.Map[String, String]].toMap } val secrets = getSecretsFromFile("/vault/secrets/appsecrets.yaml") val accessKey = secrets("presto.s3.access-key") val secretKey = secrets("presto.s3.secret-key") // 创建配置并设置凭证 val config = new Configuration() config.setString("presto.s3.access-key", accessKey) config.setString("presto.s3.secret-key", secretKey) // 补充其他必要的Flink配置 config.setString("presto.s3.endpoint", "https://s3-example-api.dev.net") config.setString("presto.s3.path-style-access", "true") // 提前初始化FileSystem,确保凭证生效 val fsFactory = new PrestoS3FileSystemFactory() fsFactory.configure(config) // 创建执行环境 val env = StreamExecutionEnvironment.getExecutionEnvironment(config)
此方案确保FileSystem在初始化时就能获取到凭证,解决BlobStore无法访问S3的问题。
内容的提问来源于stack exchange,提问作者Kubus
相关产品推荐
相关产品推荐

