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

Airflow项目中读取JSON文件并发起REST请求的问题解决

修正代码与遍历请求实现方案

一、修复read_payloads函数的路径与语法问题

原函数存在几个导致报错的核心问题:

  1. 目录名称错误:实际目录是JSON/payloads,但代码里写的是./json/payload,名称不匹配
  2. 文件路径不完整:os.listdir返回的只是文件名,直接用os.path.isfile(filename)会检查当前工作目录而非目标目录,必须拼接完整路径
  3. 字典更新语法错误:误用集合{filename, json.load(f)},应该是键值对形式{filename: json.load(f)}
  4. 未返回字典:函数末尾没有return payload_dict,调用后会得到None
  5. 文件资源未安全释放:手动打开文件后未确保关闭,建议用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 18:40:30