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

Linux环境下Asyncio调用UKG API出现JSONDecodeError求助

解决方案:异步调用UKG API随机出现JSONDecodeError问题

问题根源分析

json.decoder.JSONDecodeError: Expecting value: line 1 column 1 本质是API返回的响应为空或非JSON格式,常见触发原因包括:

  • 网络波动导致请求未正常接收完整响应
  • UKG API对并发请求限流,返回空内容/错误页面而非标准JSON
  • 部分请求实际返回4xx/5xx状态码,但未被处理就直接解析JSON
  • 超时设置不够精细,请求卡住后返回不完整的无效响应

针对性修复方案

1. 先校验响应状态码,跳过无效请求

不要直接对所有响应调用.json(),先判断状态码是否正常,异常响应记录细节方便排查:

async def get_data():
    session_timeout = aiohttp.ClientTimeout(total=1800, connect_timeout=30)
    async with aiohttp.ClientSession(timeout=session_timeout) as session:
        tasks = get_tasks(session)
        responses = await asyncio.gather(*tasks)
        for idx, response in enumerate(responses):
            emp_id = emp_ids[idx]
            if response.status != 200:
                print(f"请求员工ID {emp_id} 失败,状态码: {response.status}")
                logger.error(f"请求员工ID {emp_id} 失败,状态码: {response.status},响应内容: {await response.text()}")
                continue
            try:
                data = await response.json(content_type=None)
                response_results.append(data)
            except Exception as e:
                print(f"解析员工ID {emp_id} 响应失败: {str(e)}")
                logger.error(f"解析员工ID {emp_id} 响应失败: {str(e)},响应内容: {await response.text()}")

2. 限制并发请求数,避免触发API限流

UKG API通常对单IP的并发请求有阈值,一次性发起所有请求容易被限流。用asyncio.Semaphore控制并发数:

# 重构请求逻辑,配合Semaphore控制并发
async def fetch_emp_data(session, emp_id, semaphore):
    async with semaphore:
        async with session.get(baseURL+f"personnel/v1/user-defined-fields?employeeId={emp_id}", headers=headers) as response:
            return emp_id, response

async def get_data():
    # 限制并发数为10,可根据API文档或测试结果调整
    semaphore = asyncio.Semaphore(10)
    session_timeout = aiohttp.ClientTimeout(total=1800, connect_timeout=30, sock_read_timeout=60)
    async with aiohttp.ClientSession(timeout=session_timeout) as session:
        tasks = [fetch_emp_data(session, emp_id, semaphore) for emp_id in emp_ids]
        results = await asyncio.gather(*tasks)
        for emp_id, response in results:
            # 后续状态检查、解析逻辑同方案1
            if response.status != 200:
                print(f"请求员工ID {emp_id} 失败,状态码: {response.status}")
                logger.error(f"请求员工ID {emp_id} 失败,状态码: {response.status},响应内容: {await response.text()}")
                continue
            try:
                data = await response.json(content_type=None)
                response_results.append(data)
            except Exception as e:
                print(f"解析员工ID {emp_id} 响应失败: {str(e)}")
                logger.error(f"解析员工ID {emp_id} 响应失败: {str(e)},响应内容: {await response.text()}")

3. 优化超时配置,避免请求无响应

细化超时参数,避免因单个请求卡住导致整个批次出问题:

session_timeout = aiohttp.ClientTimeout(
    total=1800,  # 整个会话的总超时
    connect_timeout=30,  # 连接API服务器的超时
    sock_read_timeout=60,  # 读取响应内容的超时
    sock_connect_timeout=30  # 套接字层面的连接超时
)

4. 排查服务器网络稳定性

在Red Hat服务器上用curl测试UKG API的连通性,模拟请求看是否偶尔返回空:

curl -v "你的UKG API完整地址?employeeId=测试用员工ID" -H "Authorization: 你的认证token" -H "Content-Type: application/json"

如果频繁出现空响应,联系运维排查服务器到UKG API的网络链路是否有丢包、延迟过高的情况。

修改后完整代码示例

import asyncio
import aiohttp
import time
import logging

# 初始化日志
logger = logging.getLogger(__name__)
logging.basicConfig(level=logging.INFO)

# Define Global Variables
emp_ids = []  # 填充你的员工ID列表
response_results = []
baseURL = "你的UKG API基础地址"
headers = {}  # 填充你的请求头(如Authorization等)

async def fetch_emp_data(session, emp_id, semaphore):
    async with semaphore:
        try:
            async with session.get(
                baseURL + f"personnel/v1/user-defined-fields?employeeId={emp_id}",
                headers=headers
            ) as response:
                return emp_id, response
        except Exception as e:
            logger.error(f"请求员工ID {emp_id} 时发生网络错误: {str(e)}")
            return emp_id, None

async def get_data():
    # 限制并发数,可根据实际情况调整
    semaphore = asyncio.Semaphore(10)
    session_timeout = aiohttp.ClientTimeout(
        total=1800,
        connect_timeout=30,
        sock_read_timeout=60,
        sock_connect_timeout=30
    )
    async with aiohttp.ClientSession(timeout=session_timeout) as session:
        tasks = [fetch_emp_data(session, emp_id, semaphore) for emp_id in emp_ids]
        results = await asyncio.gather(*tasks)
        for emp_id, response in results:
            if response is None:
                continue
            if response.status != 200:
                resp_text = await response.text()
                logger.error(f"员工ID {emp_id} 请求失败,状态码: {response.status},响应内容: {resp_text[:200]}")
                continue
            try:
                data = await response.json(content_type=None)
                response_results.append(data)
            except Exception as e:
                resp_text = await response.text()
                logger.error(f"员工ID {emp_id} 响应解析失败: {str(e)},响应内容: {resp_text[:200]}")

if __name__ == "__main__":
    # 假设这里有获取初始连接验证response的逻辑
    # response = ... 
    if response.status_code == 200:
        print(f"Successfully connected to UKG API: Status Code {response.status_code}")
        logger.info(f'Successfully connected to UKG API: Status Code {response.status_code}')

        start_time = time.time()
        asyncio.run(get_data())
        end_time = time.time()
        total_time = end_time - start_time
        print(f"It took {total_time:.2f} seconds to make {len(emp_ids)} API Calls, successfully got {len(response_results)} results")
        logger.info(f"It took {total_time:.2f} seconds to make {len(emp_ids)} API Calls, successfully got {len(response_results)} results")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 23:18:13