如何为K8s中FlinkSessionJob传递Kafka及数据库凭证?
针对你在Flink Operator 1.10版本下,会话集群模式中FlinkSessionJob无法直接配置env参数的问题,以下是几种可行的实现方案:
方案一:通过作业启动参数传递凭证
利用FlinkSessionJob的job.arguments字段传递参数,结合K8s Secret挂载实现安全传递:
在FlinkDeployment中挂载所需Secret
修改FlinkDeployment的JobManager和TaskManager的podTemplate,将存储凭证的Secret挂载为文件:apiVersion: flink.apache.org/v1beta1 kind: FlinkDeployment metadata: name: basic-cluster spec: # 原有配置不变 jobManager: podTemplate: spec: volumes: - name: kafka-secrets-job1 secret: secretName: kafka-creds-job1 - name: db-secrets-job1 secret: secretName: db-creds-job1 containers: - name: flink-main-container volumeMounts: - name: kafka-secrets-job1 mountPath: /opt/flink/secrets/kafka/job1 readOnly: true - name: db-secrets-job1 mountPath: /opt/flink/secrets/db/job1 readOnly: true taskManager: podTemplate: spec: volumes: - name: kafka-secrets-job1 secret: secretName: kafka-creds-job1 - name: db-secrets-job1 secret: secretName: db-creds-job1 containers: - name: flink-main-container volumeMounts: - name: kafka-secrets-job1 mountPath: /opt/flink/secrets/kafka/job1 readOnly: true - name: db-secrets-job1 mountPath: /opt/flink/secrets/db/job1 readOnly: true在FlinkSessionJob中指定参数引用挂载的文件
通过arguments传递Kafka主题和凭证文件路径,作业代码中读取这些路径的内容:apiVersion: flink.apache.org/v1beta1 kind: FlinkSessionJob metadata: name: flink-job-example spec: deploymentName: basic-cluster job: jarURI: https://URL_to_my_app-0.0.1.jar parallelism: 2 upgradeMode: stateless arguments: - "--kafka-topic=topic-job1" - "--kafka-username-path=/opt/flink/secrets/kafka/job1/username" - "--kafka-password-path=/opt/flink/secrets/kafka/job1/password" - "--db-url=jdbc:mysql://db-host:3306/db-job1" - "--db-username-path=/opt/flink/secrets/db/job1/username" - "--db-password-path=/opt/flink/secrets/db/job1/password"
方案二:为每个作业挂载独立配置Secret并指定配置文件
将每个作业的完整配置(Kafka主题、数据库凭证)存入独立Secret,挂载到Session集群后,通过作业参数指定配置文件路径:
创建作业专属配置Secret
比如为job1创建包含config.properties的Secret:cat << EOF > config-job1.properties kafka.topic=topic-job1 kafka.username=xxx kafka.password=yyy db.url=jdbc:mysql://db-host:3306/db-job1 db.username=aaa db.password=bbb EOF kubectl create secret generic job1-config --from-file=config-job1.properties在FlinkDeployment中挂载该Secret
同方案一,在JobManager和TaskManager的podTemplate中添加volume和volumeMount:# FlinkDeployment的jobManager.podTemplate.spec部分 volumes: - name: job1-config secret: secretName: job1-config containers: - name: flink-main-container volumeMounts: - name: job1-config mountPath: /opt/flink/config/job1 readOnly: true # TaskManager部分同样添加上述配置FlinkSessionJob指定配置文件路径
apiVersion: flink.apache.org/v1beta1 kind: FlinkSessionJob metadata: name: flink-job-example spec: deploymentName: basic-cluster job: jarURI: https://URL_to_my_app-0.0.1.jar parallelism: 2 upgradeMode: stateless arguments: - "--config-path=/opt/flink/config/job1/config-job1.properties"
方案三:利用FlinkSessionJob的Flink配置扩展
Flink Operator 1.10版本支持在FlinkSessionJob的job.flinkConfiguration字段中添加作业级配置项,可结合Secret环境变量实现安全传递:
在FlinkDeployment中挂载Secret为环境变量
# FlinkDeployment的jobManager.podTemplate.spec.containers部分 env: - name: JOB1_KAFKA_USERNAME valueFrom: secretKeyRef: name: kafka-creds-job1 key: username - name: JOB1_KAFKA_PASSWORD valueFrom: secretKeyRef: name: kafka-creds-job1 key: password - name: JOB1_DB_USERNAME valueFrom: secretKeyRef: name: db-creds-job1 key: username - name: JOB1_DB_PASSWORD valueFrom: secretKeyRef: name: db-creds-job1 key: password # TaskManager的podTemplate同样添加上述环境变量配置在FlinkSessionJob中引用环境变量作为配置项
apiVersion: flink.apache.org/v1beta1 kind: FlinkSessionJob metadata: name: flink-job-example spec: deploymentName: basic-cluster job: jarURI: https://URL_to_my_app-0.0.1.jar parallelism: 2 upgradeMode: stateless flinkConfiguration: job.kafka.topic: "topic-job1" job.kafka.username: "$(JOB1_KAFKA_USERNAME)" job.kafka.password: "$(JOB1_KAFKA_PASSWORD)" job.db.url: "jdbc:mysql://db-host:3306/db-job1" job.db.username: "$(JOB1_DB_USERNAME)" job.db.password: "$(JOB1_DB_PASSWORD)"
内容的提问来源于stack exchange,提问作者Алексей Коротков

