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

如何将基于线程的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)

关键改进说明

  1. 异步HTTP库替换:用aiohttp替代requests,因为requests是同步阻塞的,无法在asyncio环境中发挥异步优势;
  2. 异步生成器实现:get_facts改为异步生成器(async def + yield),每获取到一个fact就立即输出,不用等待所有请求完成;
  3. 即时处理逻辑:每个fact就绪后,立刻发起和上一个fact的相似度请求,不用等所有fact拉取完成,最大化并发效率;
  4. 状态管理优化:把原来的类共享属性(facts、percentages等)改为实例属性,避免多实例运行时的状态冲突;
  5. 并发任务管理:用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 00:45:55