Airflow项目中读取JSON文件并发起REST请求的问题解决
修正代码与遍历请求实现方案
一、修复read_payloads函数的路径与语法问题
原函数存在几个导致报错的核心问题:
- 目录名称错误:实际目录是
JSON/payloads,但代码里写的是./json/payload,名称不匹配 - 文件路径不完整:
os.listdir返回的只是文件名,直接用os.path.isfile(filename)会检查当前工作目录而非目标目录,必须拼接完整路径 - 字典更新语法错误:误用集合
{filename, json.load(f)},应该是键值对形式{filename: json.load(f)} - 未返回字典:函数末尾没有
return payload_dict,调用后会得到None - 文件资源未安全释放:手动打开文件后未确保关闭,建议用
with语句自动管理
修复后的函数:
import os import json def read_payloads(category): payload_dict = {} # 用os.path.join拼接路径,适配不同系统,同时修正目录名称 directory = os.path.join("./JSON", "payloads", category) # 先校验目录是否存在,避免直接报错 if not os.path.isdir(directory): print(f"目录不存在: {directory}") return payload_dict for filename in os.listdir(directory): file_path = os.path.join(directory, filename) # 只处理文件,跳过子目录 if os.path.isfile(file_path): # with语句自动关闭文件,无需手动调用close() with open(file_path, 'r', encoding='utf-8') as f: try: payload = json.load(f) payload_dict[filename] = payload except json.JSONDecodeError: print(f"文件{filename}不是有效的JSON格式,跳过") return payload_dict
二、遍历字典发起REST请求
在server_call函数中,遍历payload_dict的键值对,逐个发送请求。这里用requests库示例(Airflow环境一般已预装):
import requests def server_call(url, token, category): payload_dict = read_payloads(category) if not payload_dict: print("没有可处理的payload数据") return # 构建请求头,根据你的API要求调整 headers = { "Authorization": f"Bearer {token}", "Content-Type": "application/json" } # 遍历每个文件名对应的payload for filename, payload in payload_dict.items(): try: # 假设是POST请求,根据API实际方法调整(GET/PUT等) response = requests.post(url, json=payload, headers=headers) response.raise_for_status() # 捕获HTTP状态码错误(如4xx/5xx) print(f"文件[{filename}]请求成功,响应状态码: {response.status_code}") # 可按需处理响应内容,比如打印或存储response.json() except requests.exceptions.RequestException as e: print(f"文件[{filename}]请求失败: {str(e)}")
三、额外注意事项
- Airflow路径适配:如果是在Airflow DAG中使用,建议使用绝对路径(比如
/opt/airflow/dags/JSON/payloads),避免因Airflow的工作目录导致路径错误 - 异常扩展:可以根据业务需求,增加更多异常处理(比如权限错误、网络超时等)
- 日志记录:在Airflow中推荐使用
logging模块替代print,方便在日志中查看执行情况
内容的提问来源于stack exchange,提问作者azaveri7
相关产品推荐
相关产品推荐

