如何实时获取多组可变参数API的JSON数据且不遗漏更新?
实时拉取多份JSON数据解决方案
核心适配结论
不用更换开发语言,你当前的Python技术栈完全可以满足需求,仅需替换原有的同步轮询逻辑为异步并发架构即可解决漏采、卡顿问题。
核心优化思路
原有的串行同步轮询存在两个核心问题:一是单线程依次请求100个接口的总耗时太长,导致单接口轮询间隔远大于数据最小更新周期;二是同步请求、同步入库的逻辑会互相阻塞,进一步放大轮询间隔。优化方向为:
- 单接口轮询间隔设置为小于最小数据更新周期(你的场景最小更新为0.025s,可设置为0.02s),从根源上避免漏采
- 所有接口的轮询、数据预处理、入库逻辑全部异步化,互不阻塞
- 增加数据变更校验,避免重复写入无效数据
工具推荐&实现方案
- 异步HTTP请求:用
aiohttp替换原有的requests库,支持单进程并发发起上百个HTTP请求无阻塞 - 异步任务调度:每个boll参数对应一个独立的异步协程,各自维护轮询周期,互不影响
- 异步数据库操作:用
aiomysql替换同步MySQL连接库,入库操作不会阻塞请求逻辑,数据量较大时可增加批量写入缓冲进一步提升性能 - 变更校验:每次拉取到数据后对JSON内容排序后计算哈希,和上一次存储的哈希对比,仅当哈希变化时才执行入库逻辑
最简示例代码
import aiohttp import asyncio import aiomysql import json from typing import Dict, Any # 业务配置 BOLL_PARAMS: list[str] = ["xxx1", "xxx2", ...] # 替换为你的50~100个boll参数 SINGLE_POLL_INTERVAL: float = 0.02 # 单接口轮询间隔20ms,小于最小更新周期25ms LAST_DATA_HASH: Dict[str, int] = {} async def poll_single_boll(session: aiohttp.ClientSession, boll: str, db_pool: aiomysql.Pool) -> None: """单个boll参数的独立轮询逻辑""" url = f"https://abcd.xyz/snapshot?boll={boll}" try: async with session.get(url, timeout=0.5) as resp: if resp.status != 200: return data: Dict[str, Any] = await resp.json() # 计算数据哈希判断是否变更 current_hash = hash(json.dumps(data, sort_keys=True)) if LAST_DATA_HASH.get(boll) == current_hash: return LAST_DATA_HASH[boll] = current_hash # 异步写入MySQL async with db_pool.acquire() as conn: async with conn.cursor() as cur: await cur.execute( "INSERT INTO your_table (boll_param, data_content, update_time) VALUES (%s, %s, NOW(3))", (boll, json.dumps(data)) ) await conn.commit() except Exception: # 可按需增加重试、日志记录逻辑 pass async def run_poll_tasks() -> None: # 初始化MySQL异步连接池 db_pool = await aiomysql.create_pool( host="127.0.0.1", port=3306, user="your_db_user", password="your_db_pass", db="your_db_name", minsize=5, maxsize=20 ) # 初始化全局HTTP会话 async with aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=10)) as session: # 为每个boll参数创建独立轮询协程 poll_tasks = [] for boll in BOLL_PARAMS: async def single_task(b: str): while True: await poll_single_boll(session, b, db_pool) await asyncio.sleep(SINGLE_POLL_INTERVAL) poll_tasks.append(asyncio.create_task(single_task(boll))) await asyncio.gather(*poll_tasks) # 资源释放 db_pool.close() await db_pool.wait_closed() if __name__ == "__main__": asyncio.run(run_poll_tasks())
可选优化项
- 若接口存在限流规则,可给每个协程的轮询间隔增加1~10ms的随机偏移,避免所有请求同时发起触发限流
- 可引入
tenacity库为请求逻辑增加重试机制,网络抖动时自动重试避免漏采 - 数据写入量较大时,可先将变更数据写入本地内存队列或Redis队列,用独立的协程/进程批量入库,进一步降低轮询逻辑的阻塞概率
内容的提问来源于stack exchange,提问作者ShanN
相关产品推荐
相关产品推荐

