在Airflow中使用Python解析复杂JSON并按条件取值
在Airflow中读取Admin Variables的JSON数据并按task_name提取值
步骤1:读取并解析Airflow Variables中的JSON数据
先在Airflow Admin的Variables页面创建变量(比如命名为task_config),将你的JSON配置作为变量值保存。然后通过以下代码读取并解析为Python可操作的字典:
from airflow.models import Variable import json # 读取变量的字符串值 config_str = Variable.get("task_config") # 解析字符串为Python字典 config_data = json.loads(config_str)
步骤2:按task_name筛选并提取目标值
遍历配置列表,匹配指定task_name后提取对应字段赋值给变量,示例代码如下:
# 设定要匹配的任务名称 target_task = "reading" # 可替换为"singing"或其他任务名 # 遍历配置项 for task_item in config_data["configuration"]: if task_item["task_name"] == target_task: # 根据任务类型提取对应字段 if target_task == "reading": novel_name = task_item["novel"]["name"] novel_author = task_item["novel"]["author"] story_name = task_item["story"]["name"] story_author = task_item["story"]["author"] # 验证提取结果(可选) print(f"小说:{novel_name},作者:{novel_author}") print(f"故事:{story_name},作者:{story_author}") elif target_task == "singing": pop_name = task_item["pop"]["name"] pop_author = task_item["pop"]["author"] jazz_name = task_item["jazz"]["name"] jazz_author = task_item["jazz"]["author"] print(f"流行曲目:{pop_name},作者:{pop_author}") print(f"爵士曲目:{jazz_name},作者:{jazz_author}") break # 找到匹配项后终止遍历 else: # 未找到对应task_name时的提示 print(f"未找到task_name为{target_task}的配置")
额外提示
- 这段逻辑可以直接嵌入Airflow的
PythonOperator任务函数中使用。 - 若配置结构调整,只需修改提取字段的层级路径即可适配。
- 可添加异常处理(比如
json.JSONDecodeError),避免因JSON格式错误导致任务失败。
内容的提问来源于stack exchange,提问作者Joe1988
相关产品推荐
相关产品推荐

