为何处理百万级URL时asyncio脚本执行速度大幅下降?
问题描述
我写了一个批量获取URL状态码的异步脚本,处理10000条URL时,速度大概是8分钟/万条;但处理百万级以上URL时,速度降到了20分钟/万条。这是异步编程的特性吗?该怎么优化?
附上脚本代码:
import logging import asyncio import aiohttp import numpy as np from tqdm import tqdm headers = {} # 假设此处已定义请求头 def config_logger(name, file): logger = logging.getLogger(name) logger.setLevel(logging.INFO) handler = logging.FileHandler(file) logger.addHandler(handler) return logger, handler def create_two_loggers(directory, file): status_codes, status_codes_handler = config_logger("status_codes", f"{directory}/status_codes.csv") valid_status_codes, valid_status_codes_handler = config_logger("valid_status_codes", f"{directory}/{file}") return status_codes, valid_status_codes async def get(url, sem, session, status_codes, valid_status_codes): try: async with sem, session.head(url=url, timeout=20, raise_for_status=True) as r: status = r.status if r.status <= 308: status_codes.info(f"{url},{status}") valid_status_codes.info(url) if r.status > 308: status_codes.error(f"{url},{status}") except Exception: status_codes.error(f"{url},{np.nan}") async def main(directory, file, urls): status_codes, valid_status_codes = create_two_loggers(directory, file) sem = asyncio.BoundedSemaphore(50) async with aiohttp.ClientSession(connector=aiohttp.TCPConnector(limit=200), headers=headers) as session: await tqdm.gather(*[get(url, sem, session, status_codes, valid_status_codes) for url in urls]) def fetch_status_codes(directory, file, urls): asyncio.run(main(directory, file, urls)) if __name__ == "__main__": fetch_status_codes()
问题分析与优化方案
这不是异步编程的固有特性,百万级URL处理速度下降是当前脚本的实现瓶颈导致的,主要集中在日志IO阻塞、任务调度开销、并发参数不匹配等方面,以下是具体优化方案:
1. 优化日志写入(核心瓶颈)
当前脚本用logging.FileHandler同步写入日志,高并发场景下,大量异步任务会因等待文件IO阻塞事件循环,这是速度下降的主要原因。
优化方案:
- 批量缓存日志:先将日志内容缓存到内存队列,再用单独线程批量写入文件,避免每个任务都触发同步IO。
- 使用异步日志库:比如
aiologger,支持异步写入,不会阻塞事件循环。
示例修改(批量缓存日志):
import queue import threading def log_writer(q, file_path): with open(file_path, "a", encoding="utf-8") as f: while True: line = q.get() if line is None: # 终止信号 break f.write(line + "\n") f.flush() def create_log_writer_thread(file_path): q = queue.Queue(maxsize=10000) thread = threading.Thread(target=log_writer, args=(q, file_path), daemon=True) thread.start() return q # 修改main函数初始化日志队列 async def main(directory, file, urls): status_log_q = create_log_writer_thread(f"{directory}/status_codes.csv") valid_log_q = create_log_writer_thread(f"{directory}/{file}") sem = asyncio.BoundedSemaphore(150) async with aiohttp.ClientSession(connector=aiohttp.TCPConnector(limit=200), headers=headers) as session: await tqdm.gather(*[get(url, sem, session, status_log_q, valid_log_q) for url in urls]) # 发送终止信号 status_log_q.put(None) valid_log_q.put(None) # 修改get函数的日志写入逻辑 async def get(url, sem, session, status_log_q, valid_log_q): try: async with sem, session.head(url=url, timeout=10, raise_for_status=True) as r: status = r.status if status <= 308: status_log_q.put(f"{url},{status}") valid_log_q.put(url) else: status_log_q.put(f"{url},{status}") except Exception as e: status_log_q.put(f"{url},{np.nan}")
2. 调整并发控制参数
当前BoundedSemaphore(50)限制同时运行的任务数为50,而TCPConnector(limit=200)允许的连接数是200,两者不匹配导致连接池资源浪费,百万级任务下过低的并发数会拉低整体效率。
优化方案:
- 根据网络环境调整semaphore值,比如设置为100-200(注意不要过高,避免触发目标服务器反爬或本地网络瓶颈)。
- 保持semaphore值与TCPConnector的limit值大致匹配,或让limit略大于semaphore值。
示例修改:
sem = asyncio.BoundedSemaphore(150) async with aiohttp.ClientSession(connector=aiohttp.TCPConnector(limit=200, keepalive_timeout=30), headers=headers) as session: ...
3. 避免一次性生成所有任务
await tqdm.gather(*[get(...) for url in urls])会一次性生成百万级任务对象,占用大量内存,同时给事件循环带来巨大调度开销。
优化方案:
- 分批处理URL,比如每次处理1000条,循环提交任务,降低内存占用和调度压力。
示例修改:
async def main(directory, file, urls): status_log_q = create_log_writer_thread(f"{directory}/status_codes.csv") valid_log_q = create_log_writer_thread(f"{directory}/{file}") sem = asyncio.BoundedSemaphore(150) batch_size = 1000 total = len(urls) async with aiohttp.ClientSession(connector=aiohttp.TCPConnector(limit=200), headers=headers) as session: with tqdm(total=total) as pbar: for i in range(0, total, batch_size): batch = urls[i:i+batch_size] tasks = [get(url, sem, session, status_log_q, valid_log_q) for url in batch] await asyncio.gather(*tasks) pbar.update(len(batch)) status_log_q.put(None) valid_log_q.put(None)
4. 优化异常处理与请求参数
- 细分异常类型:不要用
except Exception捕获所有异常,只捕获aiohttp相关请求异常(如ClientError、TimeoutError),减少不必要的处理开销。 - 拆分超时设置:将20秒统一超时拆分为连接超时和读取超时,避免慢请求拖慢整体进度。
- 启用HTTP/2:如果目标服务器支持,在
TCPConnector中设置force_close=False并启用HTTP/2,减少TCP握手开销。
示例修改:
from aiohttp import ClientError async def get(url, sem, session, status_log_q, valid_log_q): try: # 拆分超时为连接超时和读取超时 timeout = aiohttp.ClientTimeout(total=10, connect=5) async with sem, session.head(url=url, timeout=timeout, raise_for_status=True) as r: status = r.status if status <= 308: status_log_q.put(f"{url},{status}") valid_log_q.put(url) else: status_log_q.put(f"{url},{status}") except ClientError: status_log_q.put(f"{url},{np.nan}") except asyncio.TimeoutError: status_log_q.put(f"{url},timeout")
5. 内存优化(针对百万级URL)
如果URL从文件读取,不要一次性加载所有URL到内存,改为边读边处理,减少内存占用:
示例修改:
def read_urls_batch(file_path, batch_size=1000): with open(file_path, "r", encoding="utf-8") as f: batch = [] for line in f: url = line.strip() if url: batch.append(url) if len(batch) == batch_size: yield batch batch = [] if batch: yield batch async def main(directory, file, url_file_path): status_log_q = create_log_writer_thread(f"{directory}/status_codes.csv") valid_log_q = create_log_writer_thread(f"{directory}/{file}") sem = asyncio.BoundedSemaphore(150) total = sum(1 for _ in open(url_file_path, "r", encoding="utf-8") if _.strip()) async with aiohttp.ClientSession(connector=aiohttp.TCPConnector(limit=200), headers=headers) as session: with tqdm(total=total) as pbar: for batch in read_urls_batch(url_file_path): tasks = [get(url, sem, session, status_log_q, valid_log_q) for url in batch] await asyncio.gather(*tasks) pbar.update(len(batch)) status_log_q.put(None) valid_log_q.put(None)
内容的提问来源于stack exchange,提问作者ariyasas94
相关产品推荐
相关产品推荐

