如何将基于线程的API调用代码改写为高效的asyncio实现?
问题描述
我写了一个脚本,功能能正常运行,但实现方式不符合需求。以下是等效示例代码(仅用于演示核心逻辑),该脚本从API拉取数据,处理后提交至另一API,目前用线程实现并发API调用:
import requests from threading import Thread class WithThreads: facts = [] percentages = [] threads = [] HEADERS = { 'X-RapidAPI-Key': 'API-KEY', 'X-RapidAPI-Host': 'text-similarity-calculator.p.rapidapi.com' } @staticmethod def join_all(threads): """ Join all threads in the array and clear the array. """ while threads: thread = threads.pop() thread.join() def get_facts(self): """ Calls the first API n times with different parameters. (In real life.) Should ideally yield. """ for _ in range(10): thread = Thread(target=self.get_fact) self.threads.append(thread) thread.start() def get_fact(self): """ Calls the first API once. """ # May take long. self.facts.append(requests.get('https://catfact.ninja/fact', headers=self.HEADERS).json()['fact']) def get_percentage(self, ftext, stext): """ Calls the second API once. """ # May take long. response = requests.get('https://text-similarity-calculator.p.rapidapi.com/' f'stringcalculator.php?ftext={ftext}&stext={stext}', headers=self.HEADERS) self.percentages.append(response.json()['percentage'] + '%') def get_percentages(self): """ This is the main class to be called. Feeds the result of the first API into the second API, uses threads instead of asyncio. """ self.get_facts() self.join_all(self.threads) first = self.facts[0] previous = first for fact in self.facts[1:]: thread = Thread(target=self.get_percentage, args=(previous, fact)) self.threads.append(thread) thread.start() previous = fact thread = Thread(target=self.get_percentage, args=(previous, first)) self.threads.append(thread) thread.start() self.join_all(self.threads) return self.percentages
当前实现存在以下问题:
- 应改用asyncio实现异步并发,但重构时要么代码过于复杂,要么方法无法正常调用;
- 从API拉取的事实数据应以yield方式输出,线程难以实现这一点,而asyncio应支持该特性;
- 真实场景中,每个第一个API的返回结果就绪后,应立即发起对应的第二个API调用,asyncio应支持这种及时处理逻辑。
请问如何将该代码有效改写为asyncio实现?
解决方案
下面是基于asyncio的改写版本,完全解决你的三个问题,同时优化了状态管理和性能:
import asyncio import aiohttp class WithAsyncIO: def __init__(self): # 将类属性改为实例属性,避免多实例冲突 self.HEADERS = { 'X-RapidAPI-Key': 'API-KEY', 'X-RapidAPI-Host': 'text-similarity-calculator.p.rapidapi.com' } self.similarity_tasks = [] self.first_fact = None self.previous_fact = None async def get_fact(self, session): """异步获取单个cat fact""" async with session.get('https://catfact.ninja/fact', headers=self.HEADERS) as response: data = await response.json() return data['fact'] async def get_facts(self, session): """异步生成器:每获取到一个fact就立即yield""" for _ in range(10): # 发起异步请求,不用等待完成就继续下一个 fact = await self.get_fact(session) # 处理即时逻辑:拿到fact后立刻发起相似度请求(除了第一个) if self.first_fact is None: self.first_fact = fact self.previous_fact = fact else: # 把相似度请求加入任务列表,后台运行 task = asyncio.create_task(self.get_percentage(session, self.previous_fact, fact)) self.similarity_tasks.append(task) self.previous_fact = fact yield fact async def get_percentage(self, session, ftext, stext): """异步请求相似度API""" url = f'https://text-similarity-calculator.p.rapidapi.com/stringcalculator.php?ftext={ftext}&stext={stext}' async with session.get(url, headers=self.HEADERS) as response: data = await response.json() return f"{data['percentage']}%" async def get_percentages(self): """主方法:协调所有异步任务""" # 创建aiohttp会话,复用连接提升效率 async with aiohttp.ClientSession() as session: # 遍历异步生成器获取所有fact async for _ in self.get_facts(session): pass # 最后处理首尾fact的相似度 if self.first_fact and self.previous_fact: final_task = asyncio.create_task(self.get_percentage(session, self.previous_fact, self.first_fact)) self.similarity_tasks.append(final_task) # 等待所有相似度任务完成,收集结果 percentages = await asyncio.gather(*self.similarity_tasks) return percentages # 运行示例 if __name__ == "__main__": result = asyncio.run(WithAsyncIO().get_percentages()) print(result)
关键改进说明
- 异步HTTP库替换:用
aiohttp替代requests,因为requests是同步阻塞的,无法在asyncio环境中发挥异步优势; - 异步生成器实现:
get_facts改为异步生成器(async def+yield),每获取到一个fact就立即输出,不用等待所有请求完成; - 即时处理逻辑:每个fact就绪后,立刻发起和上一个fact的相似度请求,不用等所有fact拉取完成,最大化并发效率;
- 状态管理优化:把原来的类共享属性(
facts、percentages等)改为实例属性,避免多实例运行时的状态冲突; - 并发任务管理:用
asyncio.create_task创建后台任务,用asyncio.gather统一等待所有任务完成并收集结果。
注意事项
- 先安装依赖:
pip install aiohttp; - 替换代码中的
API-KEY为你的真实RapidAPI密钥; - 如果API有并发请求限制,可以用
asyncio.Semaphore控制并发数,示例如下:def __init__(self): # ...其他初始化代码 self.semaphore = asyncio.Semaphore(5) # 限制同时最多5个请求 async def get_fact(self, session): async with self.semaphore: # 原请求代码
内容的提问来源于stack exchange,提问作者eje211
相关产品推荐
相关产品推荐

