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

Airflow的DatabricksSubmitRunOperator如何结合Jinja模板使用?

问题根因说明

Airflow 默认的Jinja模板渲染逻辑统一返回字符串类型,即便DatabricksSubmitRunOperator的json字段属于template_field,渲染后的结果也是JSON格式的字符串,无法匹配算子要求的dict类型校验规则。

方案1:新增Jinja from_json 过滤器(全局通用)

通过自定义Jinja过滤器将渲染后的JSON字符串直接转成dict,适合多DAG复用的场景。

  1. 注册过滤器(可放在Airflow插件目录或DAG文件头部)
import json
from airflow.plugins_manager import AirflowPlugin

def from_json(json_str):
    return json.loads(json_str)

class JinjaExtendPlugin(AirflowPlugin):
    name = "jinja_extend"
    jinja_filters = {
        "from_json": from_json
    }
  1. 在DAG的算子参数中直接调用过滤器
DatabricksSubmitRunOperator(
    task_id="databricks_job_submit",
    json="{{ databricks_config_template | from_json }}",
    databricks_conn_id="your_databricks_conn",
    # 其余参数省略
)

方案2:自定义继承算子自动处理类型(算子级通用)

重写算子的模板渲染方法,自动对json字段做类型转换,无需修改原有模板写法。

from airflow.providers.databricks.operators.databricks import DatabricksSubmitRunOperator
import json

class CustomDatabricksSubmitRunOperator(DatabricksSubmitRunOperator):
    def render_template(self, attr, content, context, jinja_env=None):
        rendered_content = super().render_template(attr, content, context, jinja_env)
        # 仅对json字段做字符串转dict处理
        if attr == "json" and isinstance(rendered_content, str):
            return json.loads(rendered_content)
        return rendered_content

后续直接使用自定义的CustomDatabricksSubmitRunOperator即可,原有模板逻辑不需要调整。

方案3:XCom传递预生成dict(无侵入变通方案)

如果不想修改全局配置或自定义算子,可通过前置Python任务生成符合要求的dict配置,存入XCom后直接给Databricks算子调用,Airflow 2.x+会自动保留XCom的原始dict类型。

  1. 新增配置生成任务
from airflow.operators.python import PythonOperator

def gen_databricks_config(**context):
    # 可在这里实现任意配置生成逻辑,包括手动渲染Jinja模板得到dict
    config = {
        "new_cluster": {
            "spark_version": "13.3.x-scala2.12",
            "node_type_id": "m5.xlarge",
            "num_workers": 2
        },
        "spark_python_task": {
            "python_file": "dbfs:/path/to/your/job.py",
            "parameters": [context["ds"], context["next_ds"]]
        }
    }
    return config

gen_config_task = PythonOperator(
    task_id="generate_databricks_config",
    python_callable=gen_databricks_config,
    provide_context=True
)
  1. Databricks算子直接读取XCom配置
DatabricksSubmitRunOperator(
    task_id="databricks_job_submit",
    json="{{ ti.xcom_pull(task_ids='generate_databricks_config') }}",
    databricks_conn_id="your_databricks_conn",
    # 其余参数省略
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 16:30:04