You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.07 10:40:53