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

