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

Python asyncio.sleep内存过高及API限流等问题咨询

问题与解决方案

问题描述

我编写了一个异步函数,该函数会创建约390个任务,每处理一定任务后会执行await asyncio.sleep(1秒)。运行后发现脚本内存占用极高,通过memory profiler分析,观察到asyncio.sleep是主要内存消耗点。尝试用time.sleep替代后触发了调用API的每秒10次限流警告;改用线程池后,又遇到requests的response.json()持续占用4-8MB内存的问题,请问该如何解决?

性能分析器输出

Line #    Mem usage    Increment  Occurrences   Line Contents
=============================================================
    91    667.5 MiB    667.5 MiB           1       @profile
    92                                             async def main(self):
    93    667.5 MiB      0.0 MiB           1           tasks = []
    94    667.5 MiB      0.0 MiB           1           start_time = time.time()
    95    667.5 MiB      0.0 MiB           1           following_listings_count = len(self.following_listings)
    96
    97   1018.8 MiB      0.0 MiB           3           async with self.api.session.create_async_oauth2_session() as session:
    98
    99    988.2 MiB      0.0 MiB         351               for offset in range(0, following_listings_count, 100):
   100    988.2 MiB      0.0 MiB         350                   paged_listings = self.following_listings[offset:offset + 100]
   101    988.2 MiB      0.0 MiB       36007                   listing_ids = [listing.listing_id for listing in paged_listings]
   102
   103    988.2 MiB      0.0 MiB         700                   tasks.append(
   104    988.2 MiB      0.0 MiB         350                       asyncio.create_task(self.get_listings(session, listing_ids, paged_listings))
   105                                                         )
   106    988.2 MiB      0.0 MiB         350                   if offset % 1000 == 0:
   107    988.2 MiB    320.6 MiB          70                      await asyncio.sleep(1.1)
   108    988.2 MiB      0.0 MiB         350                   print(offset)
   109   1018.8 MiB     30.6 MiB           2               data = await asyncio.gather(*tasks)
   110
   111
   112   1018.8 MiB      0.0 MiB           1               print(f"Data Scraping Finish Time: {time.time() - start_time}")
   113   1018.8 MiB      0.0 MiB           1           """start_time = time.time()
   114                                                 all_listings = []
   115                                                 for listings in data:
   116                                                     all_listings += listings
   117
   118                                                 self.create_reports(all_listings)
   119
   120                                                 print(f"Finished Time: {time.time() - start_time}")
   121                                                 """
   122   1018.8 MiB      0.0 MiB           1           return "Ok."

解决方案

解决asyncio.sleep内存过高问题

  • 分批执行并清理任务:当前代码把所有任务一次性加入列表,直到最后才await asyncio.gather,导致数百个未完成任务堆积在内存中。改为每批执行完就清理任务,避免资源滞留:
    async with self.api.session.create_async_oauth2_session() as session:
        tasks = []
        for offset in range(0, following_listings_count, 100):
            paged_listings = self.following_listings[offset:offset + 100]
            listing_ids = [listing.listing_id for listing in paged_listings]
            task = asyncio.create_task(self.get_listings(session, listing_ids, paged_listings))
            tasks.append(task)
            
            # 每处理1000条或最后一批时,执行任务并清理
            if offset % 1000 == 0 or offset + 100 >= following_listings_count:
                await asyncio.gather(*tasks)
                tasks.clear()  # 释放已完成任务的内存
                await asyncio.sleep(1.1)
            print(offset)
    
  • 检查get_listings的资源释放:asyncio.sleep被标记为内存消耗点,实际可能是任务持有未释放的大对象(如API响应)。确保get_listings处理完数据后,及时将大对象赋值为None,或用生成器逐步处理,避免一次性缓存所有数据。

解决response.json()内存占用问题

  • 流式解析JSON:用ijson库替代response.json(),只提取需要的字段,避免加载整个JSON到内存:
    import ijson
    
    # 线程池中的requests请求示例
    def get_listings(...):
        with requests.get(url) as response:
            parser = ijson.parse(response.raw)
            for prefix, event, value in parser:
                # 仅处理所需字段,比如提取listing_id和关键数据
                if prefix.endswith('.listing_id'):
                    # 处理数据,不保存完整JSON
    
  • 强制释放响应资源:用with语句包裹请求,确保响应处理完后自动关闭;若未用with,手动调用response.close()释放连接和内存。

解决限流问题

  • 用信号量控制并发:替代sleep限流,用asyncio.Semaphore限制同时发起的API请求数,既避免触发限流,又不会导致任务堆积:
    # 在main函数外或类中定义信号量
    semaphore = asyncio.Semaphore(10)  # 对应API每秒10次的限制
    
    async def get_listings(session, listing_ids, paged_listings):
        async with semaphore:
            # 发起API请求、处理响应的逻辑
    

内容的提问来源于stack exchange,提问作者Esat Yılmaz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 13:15:00