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
相关产品推荐
相关产品推荐

