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

为何处理百万级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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 01:05:27