如何在FlinkDeployment的flinkConfiguration中使用环境变量?
问题描述
我正在使用Flink Kubernetes Operator将Flink部署到Kubernetes集群中,现有如下FlinkDeployment配置:
apiVersion: flink.apache.org/v1beta1 kind: FlinkDeployment metadata: name: myapp spec: flinkVersion: "v1_15" flinkConfiguration: kubernetes.entry.path: "/opt/flink/goldsky-docker-entrypoint.sh" taskmanager.numberOfTaskSlots: "2" ...otherstuff
我想要配置StatsD上报器,并希望通过环境变量配置其主机地址和端口,预期配置示例如下:
apiVersion: flink.apache.org/v1beta1 kind: FlinkDeployment metadata: name: myapp spec: flinkVersion: "v1_15" flinkConfiguration: kubernetes.entry.path: "/opt/flink/goldsky-docker-entrypoint.sh" taskmanager.numberOfTaskSlots: "2" metrics.reporter.stsd.factory.class: org.apache.flink.metrics.statsd.StatsDReporterFactory metrics.reporter.stsd.host: ENV_VAR_STATSD_HOST metrics.reporter.stsd.port: ENV_VAR_STATSD_PORT metrics.reporter.stsd.interval: 60 SECONDS
其中主机地址需要通过环境变量配置,因为它依赖于Flink Pod所在节点的IP地址(我需要将StatsD指标发送到集群中以DaemonSet形式运行的Vector)。端口可以硬编码,但我也希望通过环境变量配置。请问这种配置是否可行?如何通过Kubernetes Operator在Flink配置中使用环境变量?
解决方案
直接写环境变量名的配置不可行
Flink配置项无法直接识别环境变量名称,Flink Kubernetes Operator也不会自动将配置中的ENV_VAR_STATSD_HOST这类字符串替换为实际的环境变量值。可以通过以下两种方式实现需求:
方案一:自定义入口脚本+Downward API注入节点IP
1. 配置Pod模板注入环境变量
在FlinkDeployment的spec.podTemplate中,用Kubernetes的Downward API注入当前节点IP,同时定义端口的环境变量:
apiVersion: flink.apache.org/v1beta1 kind: FlinkDeployment metadata: name: myapp spec: flinkVersion: "v1_15" flinkConfiguration: kubernetes.entry.path: "/opt/flink/goldsky-docker-entrypoint.sh" taskmanager.numberOfTaskSlots: "2" # 先写占位符,后续脚本替换 metrics.reporter.stsd.factory.class: org.apache.flink.metrics.statsd.StatsDReporterFactory metrics.reporter.stsd.host: "${STATSD_HOST}" metrics.reporter.stsd.port: "${STATSD_PORT}" metrics.reporter.stsd.interval: 60 SECONDS podTemplate: spec: containers: - name: flink-main-container env: # 通过Downward API获取节点IP - name: STATSD_HOST valueFrom: fieldRef: fieldPath: status.hostIP # 自定义端口环境变量,也可以直接硬编码 - name: STATSD_PORT value: "8125"
2. 修改自定义入口脚本
在你的goldsky-docker-entrypoint.sh中,添加替换占位符的逻辑,在启动Flink前修改flink-conf.yaml:
#!/bin/bash # 替换flink-conf.yaml中的占位符 sed -i "s/\${STATSD_HOST}/${STATSD_HOST}/g" /opt/flink/conf/flink-conf.yaml sed -i "s/\${STATSD_PORT}/${STATSD_PORT}/g" /opt/flink/conf/flink-conf.yaml # 执行原Flink入口逻辑 exec /opt/flink/bin/standalone-job.sh start-foreground "$@"
方案二:利用Flink系统属性引用环境变量
Flink支持在配置中使用${sys:变量名}引用JVM系统属性,我们可以通过环境变量传递系统属性:
apiVersion: flink.apache.org/v1beta1 kind: FlinkDeployment metadata: name: myapp spec: flinkVersion: "v1_15" flinkConfiguration: kubernetes.entry.path: "/opt/flink/goldsky-docker-entrypoint.sh" taskmanager.numberOfTaskSlots: "2" metrics.reporter.stsd.factory.class: org.apache.flink.metrics.statsd.StatsDReporterFactory # 使用sys前缀引用系统属性 metrics.reporter.stsd.host: "${sys:statsd.host}" metrics.reporter.stsd.port: "${sys:statsd.port}" metrics.reporter.stsd.interval: 60 SECONDS # 传递系统属性到JVM env.java.opts: "-Dstatsd.host=$(STATSD_HOST) -Dstatsd.port=$(STATSD_PORT)" podTemplate: spec: containers: - name: flink-main-container env: - name: STATSD_HOST valueFrom: fieldRef: fieldPath: status.hostIP - name: STATSD_PORT value: "8125"
这种方式不需要修改入口脚本,Flink启动时会自动解析${sys:xxx}并替换为对应的JVM系统属性值。
内容的提问来源于stack exchange,提问作者Paymahn Moghadasian
相关产品推荐
相关产品推荐

