如何在不暴露凭证的情况下向DatabricksSubmitNowOperator传递Airflow连接凭证?
安全传递Airflow Postgres凭证到DatabricksSubmitNowOperator的方案
以下是几种可行的安全方案,避免凭证在Databricks任务详情中暴露:
方案1:通过Airflow将凭证写入Databricks Secrets,任务中引用密钥
- 在Databricks中创建专属的secret scope(例如
airflow-managed-secrets),可通过Databricks CLI或UI完成:databricks secrets create-scope --scope airflow-managed-secrets --initial-manage-principal users - 在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 ) - 修改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连接 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传递(需确认集群配置允许任务级环境变量):
- 在Airflow UI中创建加密的Variable(如
postgres_user、postgres_password,确保Airflow已配置加密后端)。 - 修改
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") } ) - 在Databricks任务代码中读取环境变量:
import os user = os.environ.get("POSTGRES_USER") password = os.environ.get("POSTGRES_PASSWORD")
内容的提问来源于stack exchange,提问作者qni dopi
相关产品推荐
相关产品推荐

