Airflow:如何向Databricks算子传递字典格式变量?
解决Airflow中字典变量传入DatabricksRunNowOperator的问题
方法1:利用Airflow Jinja模板直接解析字典
DatabricksRunNowOperator的参数支持Jinja模板渲染,无需提前拼接字符串,直接在notebook_params中用Jinja表达式获取字典并解包即可:
dict_name = "dict1" task1 = DatabricksRunNowOperator( task_id=f'Databricks_{dict_name}', databricks_conn_id='databricks', job_id=1111, notebook_params={ "param1": "param1", **{{ var.json.dictionaries[dict_name] }} } )
Airflow会自动将var.json.dictionaries[dict_name]解析为Python字典,不需要手动处理格式转换。
方法2:使用Airflow Variable类直接获取字典
如果更倾向于显式获取变量,可通过Variable.get方法直接拿到字典类型数据,指定deserialize_json=True即可自动反序列化:
from airflow.models import Variable # 获取顶层字典后提取目标子字典 dictionaries = Variable.get("dictionaries", deserialize_json=True) dict_values = dictionaries["dict1"] task1 = DatabricksRunNowOperator( task_id='Databricks_dict1', databricks_conn_id='databricks', job_id=1111, notebook_params={"param1": "param1", **dict_values} )
你之前的错误原因
- 用f-string生成的
values是包含Jinja语法的字符串("{{ var.json.dictionaries.dict1 }}"),而非实际解析后的字典对象,因此用**解包字符串会触发TypeError: 'str' object is not a mapping。 - 尝试
json.loads转换时,该字符串是Jinja模板语法,并非有效的JSON格式,自然会出现解码错误。
内容的提问来源于stack exchange,提问作者Óscar
相关产品推荐
相关产品推荐

