如何在Flink Kubernetes集群中为Python作业传递环境变量
问题描述
我正在使用Flink Kubernetes Operator 1.3.0,需要为Python作业传递一些环境变量。已按照官方文档运行示例作业且一切正常,请问如何注入环境变量以便在Python文件中使用?
补充信息:
使用的YAML文件(来自官方示例):
apiVersion: flink.apache.org/v1beta1 kind: FlinkDeployment metadata: name: python-example spec: image: localhost:32000/flink-python-example:1.16.0 flinkVersion: v1_16 flinkConfiguration: taskmanager.numberOfTaskSlots: "1" serviceAccount: flink jobManager: resource: memory: "2048m" cpu: 1 taskManager: resource: memory: "2048m" cpu: 1 job: jarURI: local:///opt/flink/opt/flink-python_2.12-1.16.0.jar # 该jarURI仅为占位符 entryClass: "org.apache.flink.client.python.PythonDriver" args: ["-pyclientexec", "/usr/local/bin/python3", "-py", "/opt/flink/usrlib/python_demo.py"] parallelism: 1 upgradeMode: stateless
对应的Python代码:
import logging import sys from pyflink.datastream import StreamExecutionEnvironment from pyflink.table import StreamTableEnvironment def python_demo(): env = StreamExecutionEnvironment.get_execution_environment() env.set_parallelism(1) t_env = StreamTableEnvironment.create(stream_execution_environment=env) t_env.execute_sql(""" CREATE TABLE orders ( order_number BIGINT, price DECIMAL(32,2), buyer ROW<first_name STRING, last_name STRING>, order_time TIMESTAMP(3) ) WITH ( 'connector' = 'datagen' )""") t_env.execute_sql(""" CREATE TABLE print_table WITH ('connector' = 'print') LIKE orders""") t_env.execute_sql(""" INSERT INTO print_table SELECT * FROM orders""") if __name__ == '__main__': logging.basicConfig(stream=sys.stdout, level=logging.INFO, format="%(message)s") python_demo()
解决方法
要给Flink Python作业注入环境变量,需在FlinkDeployment的jobManager和taskManager配置块中添加env字段——因为Python作业的驱动逻辑运行在JobManager,任务执行逻辑运行在TaskManager,两端都需要访问环境变量。
1. 修改YAML配置
在原有jobManager和taskManager节点下添加env配置,示例如下:
apiVersion: flink.apache.org/v1beta1 kind: FlinkDeployment metadata: name: python-example spec: image: localhost:32000/flink-python-example:1.16.0 flinkVersion: v1_16 flinkConfiguration: taskmanager.numberOfTaskSlots: "1" serviceAccount: flink jobManager: resource: memory: "2048m" cpu: 1 env: # JobManager环境变量配置 - name: CUSTOM_ENV_VAR value: "flink_python_demo" - name: DB_PASSWORD valueFrom: # 从K8s Secret读取敏感变量 secretKeyRef: name: db-secret key: password taskManager: resource: memory: "2048m" cpu: 1 env: # TaskManager环境变量配置 - name: CUSTOM_ENV_VAR value: "flink_python_demo" - name: DB_PASSWORD valueFrom: secretKeyRef: name: db-secret key: password job: jarURI: local:///opt/flink/opt/flink-python_2.12-1.16.0.jar entryClass: "org.apache.flink.client.python.PythonDriver" args: ["-pyclientexec", "/usr/local/bin/python3", "-py", "/opt/flink/usrlib/python_demo.py"] parallelism: 1 upgradeMode: stateless
配置说明:
- 若需全局统一变量,确保
jobManager和taskManager的env配置一致; - 支持直接赋值(
value)或从K8s Secret/ConfigMap读取(valueFrom),适配不同敏感程度的变量需求。
2. 在Python代码中读取变量
导入os模块,使用os.getenv()方法读取注入的环境变量:
import logging import sys import os # 新增导入os模块 from pyflink.datastream import StreamExecutionEnvironment from pyflink.table import StreamTableEnvironment def python_demo(): env = StreamExecutionEnvironment.get_execution_environment() env.set_parallelism(1) # 读取环境变量,第二个参数为变量不存在时的默认值 custom_var = os.getenv("CUSTOM_ENV_VAR", "default_value") db_pwd = os.getenv("DB_PASSWORD") logging.info(f"读取到的CUSTOM_ENV_VAR: {custom_var}") t_env = StreamTableEnvironment.create(stream_execution_environment=env) t_env.execute_sql(""" CREATE TABLE orders ( order_number BIGINT, price DECIMAL(32,2), buyer ROW<first_name STRING, last_name STRING>, order_time TIMESTAMP(3) ) WITH ( 'connector' = 'datagen' )""") t_env.execute_sql(""" CREATE TABLE print_table WITH ('connector' = 'print') LIKE orders""") t_env.execute_sql(""" INSERT INTO print_table SELECT * FROM orders""") if __name__ == '__main__': logging.basicConfig(stream=sys.stdout, level=logging.INFO, format="%(message)s") python_demo()
3. 应用配置并验证
执行命令更新部署:
kubectl apply -f your-flink-deployment.yaml
查看JobManager或TaskManager的Pod日志,确认环境变量已正确注入并能被Python代码读取。
内容的提问来源于stack exchange,提问作者sunny
相关产品推荐
相关产品推荐

