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

asyncio中独立进程运行CPU密集任务为何出现显著性能下降

问题背景

我当前使用aiohttp搭配asyncio发起异步网络请求,提取HTML页面文本后将结果保存到本地。项目中使用BeautifulSoup(封装在extract_text()方法内)处理响应内容,过滤代码等无关信息、提取HTML页面有效文本,但遇到异步+多进程版本脚本运行速度慢于纯同步版本的问题。

根据技术原理,BeautifulSoup解析属于CPU密集型操作,会在parse()执行阶段阻塞主事件循环,因此参考相关技术方案,选择将extract_text()放到独立进程中运行,避免事件循环阻塞。
但实际测试显示,该实现版本耗时比无多进程的同步版本高1.5倍。

为排查异步代码本身的实现问题,我移除了extract_text()调用,直接保存响应对象返回的原始文本,此时异步代码运行速度远高于同步版本,证明性能问题完全来自独立进程运行extract_text()的环节。


核心实现代码
import asyncio
from asyncio import Semaphore
import json
import logging
from pathlib import Path
from typing import List, Optional

import aiofiles
from aiohttp import ClientSession
import aiohttp
from bs4 import BeautifulSoup
import concurrent.futures
import functools


def extract_text(raw_text: str) -> str:
    return " ".join(BeautifulSoup(raw_text, "html.parser").stripped_strings)


async def fetch_text(
    url: str,
    session: ClientSession,
    semaphore: Semaphore,
    **kwargs: dict,
) -> str:
    async with semaphore:
        response = await session.request(method="GET", url=url, **kwargs)
        response.raise_for_status()
        logging.info("Got response [%s] for URL: %s", response.status, url)
        text = await response.text(encoding="utf-8")
        return text


async def parse(
    url: str,
    session: ClientSession,
    semaphore: Semaphore,
    **kwargs,
) -> Optional[str]:
    try:
        text = await fetch_text(
            url=url,
            session=session,
            semaphore=semaphore,
            **kwargs,
        )
    except (
        aiohttp.ClientError,
        aiohttp.http_exceptions.HttpProcessingError,
    ) as e:
        logging.error(
            "aiohttp exception for %s [%s]: %s",
            url,
            getattr(e, "status", None),
            getattr(e, "message", None),
        )
    except Exception as e:
        logging.exception(
            "Non-aiohttp exception occured:  %s",
            getattr(e, "__dict__", None),
        )
    else:
        loop = asyncio.get_running_loop()
        with concurrent.futures.ProcessPoolExecutor() as pool:
            extract_text_ = functools.partial(extract_text, text)
            text = await loop.run_in_executor(pool, extract_text_)
            logging.info("Found text for %s", url)
            return text


async def process_file(
    url: dict,
    session: ClientSession,
    semaphore: Semaphore,
    **kwargs: dict,
) -> None:
    category = url.get("category")
    link = url.get("link")
    if category and link:
        text = await parse(
            url=f"{URL}/{link}",
            session=session,
            semaphore=semaphore,
            **kwargs,
        )
        if text:
            save_path = await get_save_path(
                link=link,
                category=category,
            )
            await write_file(html_text=text, path=save_path)
        else:
            logging.warning("Text for %s not found, skipping it...", link)


async def process_files(
    html_files: List[dict],
    semaphore: Semaphore,
) -> None:
    async with ClientSession() as session:
        tasks = [
            process_file(
                url=file,
                session=session,
                semaphore=semaphore,
            )
            for file in html_files
        ]
        await asyncio.gather(*tasks)


async def write_file(
    html_text: str,
    path: Path,
) -> None:
    # 基于aiofiles实现文件写入逻辑
    ...

async def get_save_path(link: str, category: str) -> Path:
    # 返回文件存储路径
    ...

async def main_async(
    num_files: Optional[int],
    semaphore_count: int,
) -> None:
    html_files = # 获取所有待处理文件列表
    semaphore = Semaphore(semaphore_count)
    await process_files(
        html_files=html_files,
        semaphore=semaphore,
    )


if __name__ == "__main__":
    NUM_FILES = # 命令行传入参数
    SEMAPHORE_COUNT = # 命令行传入参数
    asyncio.run(
        main_async(
            num_files=NUM_FILES,
            semaphore_count=SEMAPHORE_COUNT,
        )
    )

1000次采样下的SnakeViz性能分析结果
  • 集成多进程运行extract_text的异步版本
    集成多进程运行extract_text的异步版本
  • 未调用extract_text的异步版本
    未调用extract_text的异步版本
  • 调用extract_text的同步版本(可见BeautifulSoup的html_parser占用了绝大多数运行时间)
    调用extract_text的同步版本
  • 未调用extract_text的同步版本
    未调用extract_text的同步版本

问题根因

你的实现存在3个核心缺陷,引入的额外开销完全抵消了多进程的性能收益,甚至比同步版本更慢:

  1. 进程池重复创建销毁
    你在parse()方法内通过with上下文创建ProcessPoolExecutor,相当于每处理一个URL就会新建整套进程池、任务执行完立刻销毁。进程创建、销毁的系统开销极高,这部分耗时已经超过了单进程直接执行BeautifulSoup解析的成本。
    正确做法是在程序启动时初始化一次全局进程池,worker数量设置为和CPU核心数持平,所有解析任务复用同一个池实例,程序退出前再统一回收资源。
  2. 跨进程通信开销过高
    多进程间传递参数必须经过pickle序列化+进程间通信传输,你将完整的HTML文本(单页通常几十KB到数MB)在主进程、子进程之间来回拷贝,大文本序列化和传输的开销非常高,很多场景下这部分耗时甚至超过HTML解析本身。
    如果单页HTML解析耗时低于100ms,跨进程的开销完全覆盖收益,这种场景优先替换解析器:把默认的html.parser换成lxml,解析速度可以提升3~5倍,直接用默认线程池执行解析即可,不需要启动多进程。
  3. 并发度不匹配
    你只给网络请求加了信号量控制,没有限制提交到进程池的任务数量,如果爬取并发数远大于CPU核心数,会导致进程池任务队列积压,大量任务排队等待调度,进一步拉长整体运行时间。

修正后的核心逻辑参考
import os
# 全局初始化一次进程池,worker数等于CPU核心数
PROCESS_POOL = concurrent.futures.ProcessPoolExecutor(max_workers=os.cpu_count())

async def parse(...):
    # ... 省略爬取逻辑
    loop = asyncio.get_running_loop()
    # 复用全局进程池,不要反复新建
    text = await loop.run_in_executor(PROCESS_POOL, extract_text, text)
    return text

# 程序退出时关闭进程池
async def main_async(...):
    try:
        # ... 原有业务逻辑
    finally:
        PROCESS_POOL.shutdown(wait=True)

如果换用lxml解析器后单页解析速度足够快,直接把PROCESS_POOL换成默认的线程池即可,线程间传递参数不需要序列化,额外开销低很多,足够应对CPU占用不高的解析场景。


内容的提问来源于stack exchange,提问作者yudhiesh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 11:15:29