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

如何在不暴露凭证的情况下向DatabricksSubmitNowOperator传递Airflow连接凭证?

安全传递Airflow Postgres凭证到DatabricksSubmitNowOperator的方案

以下是几种可行的安全方案,避免凭证在Databricks任务详情中暴露:

方案1:通过Airflow将凭证写入Databricks Secrets,任务中引用密钥

  1. 在Databricks中创建专属的secret scope(例如airflow-managed-secrets),可通过Databricks CLI或UI完成:
    databricks secrets create-scope --scope airflow-managed-secrets --initial-manage-principal users
    
  2. 在Airflow DAG中添加前置Python任务,用DatabricksHook调用Secrets API,将Airflow的Postgres连接凭证写入Databricks Secrets:
    from airflow.hooks.base import BaseHook
    from airflow.providers.databricks.hooks.databricks import DatabricksHook
    from airflow.operators.python import PythonOperator
    
    def write_postgres_secrets_to_databricks():
        # 获取Airflow中配置的Postgres连接
        conn = BaseHook.get_connection("postgres_conn_id")
        db_hook = DatabricksHook(databricks_conn_id="databricks_conn_id")
        
        # 写入用户名到Databricks Secrets
        db_hook.run(
            endpoint="secrets/put",
            method="POST",
            json={
                "scope": "airflow-managed-secrets",
                "key": "postgres-user",
                "string_value": conn.login
            }
        )
        # 写入密码到Databricks Secrets
        db_hook.run(
            endpoint="secrets/put",
            method="POST",
            json={
                "scope": "airflow-managed-secrets",
                "key": "postgres-password",
                "string_value": conn.password
            }
        )
    
    write_secrets_task = PythonOperator(
        task_id="write_postgres_secrets",
        python_callable=write_postgres_secrets_to_databricks
    )
    
  3. 修改Databricks任务代码,通过dbutils.secrets.get读取凭证,不再依赖base_parameters:
    user = dbutils.secrets.get("airflow-managed-secrets", "postgres-user")
    password = dbutils.secrets.get("airflow-managed-secrets", "postgres-password")
    # 后续用获取到的凭证建立Postgres连接
    
  4. DatabricksSubmitNowOperator的base_parameters仅传递非敏感参数(如数据库主机、库名),敏感凭证完全通过密钥读取。

方案2:直接调用Databricks API提交任务,用Spark配置传递机密

绕过DatabricksSubmitNowOperator,用Python直接调用Databricks Jobs API提交任务,将敏感凭证放入spark_conf中(Spark配置中的敏感值不会在Databricks任务详情中显示):

from airflow.hooks.base import BaseHook
from airflow.providers.databricks.hooks.databricks import DatabricksHook
from airflow.operators.python import PythonOperator

def submit_databricks_job_with_secure_params():
    conn = BaseHook.get_connection("postgres_conn_id")
    db_hook = DatabricksHook(databricks_conn_id="databricks_conn_id")
    
    job_payload = {
        "existing_cluster_id": "your-cluster-id",
        "spark_python_task": {
            "python_file": "dbfs:/path/to/your/script.py",
            "parameters": ["--db-host", conn.host, "--db-name", conn.schema]
        },
        "spark_conf": {
            "spark.postgres.user": conn.login,
            "spark.postgres.password": conn.password
        }
    }
    
    db_hook.run(
        endpoint="jobs/runs/submit",
        method="POST",
        json=job_payload
    )

submit_job_task = PythonOperator(
    task_id="submit_secure_databricks_job",
    python_callable=submit_databricks_job_with_secure_params
)

在Databricks任务代码中通过spark.conf.get获取凭证:

user = spark.conf.get("spark.postgres.user")
password = spark.conf.get("spark.postgres.password")

方案3:Airflow加密变量+Databricks任务级环境变量

如果现有集群支持任务级环境变量注入,可将Airflow连接的凭证存入加密的Airflow Variable,再通过environment_vars传递(需确认集群配置允许任务级环境变量):

  1. 在Airflow UI中创建加密的Variable(如postgres_user、postgres_password,确保Airflow已配置加密后端)。
  2. 修改DatabricksSubmitNowOperator参数,添加environment_vars:
    from airflow.models import Variable
    from airflow.providers.databricks.operators.databricks import DatabricksSubmitNowOperator
    
    submit_task = DatabricksSubmitNowOperator(
        task_id="submit_databricks_task",
        databricks_conn_id="databricks_conn_id",
        existing_cluster_id="your-cluster-id",
        spark_python_task={
            "python_file": "dbfs:/path/to/script.py",
            "parameters": ["--db-host", "your-db-host", "--db-name", "your-db-name"]
        },
        environment_vars={
            "POSTGRES_USER": Variable.get("postgres_user"),
            "POSTGRES_PASSWORD": Variable.get("postgres_password")
        }
    )
    
  3. 在Databricks任务代码中读取环境变量:
    import os
    user = os.environ.get("POSTGRES_USER")
    password = os.environ.get("POSTGRES_PASSWORD")
    

内容的提问来源于stack exchange,提问作者qni dopi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 16:57:47