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

Python多进程队列实现未提速问题求助

问题分析与解决方案

嘿,我帮你找到问题所在了!你的多进程代码其实根本没真正并行起来,难怪耗时和单进程差不多~

核心问题

你当前的多进程架构存在一个关键缺陷:

  • 你只启动了一个message_collector_proc进程来处理队列里的任务,所有site.all()调用都是在这个单进程里串行执行的,和单进程版本的执行逻辑几乎完全一致,自然耗时没差别。
  • 虽然你用了JoinableQueue来传递任务,但只有一个消费者进程,队列里的任务是逐个被处理的,完全没利用到多进程并行的优势。

修复方案:真正实现并行处理

我们可以用multiprocessing.Pool来简化多进程任务的管理,让每个site.all()在独立的进程中并行执行。这是Python多进程处理批量任务最常用的方式。

修改后的完整代码:

import time
import multiprocessing as mp
import requests
from bs4 import BeautifulSoup as bs

def complex_task():
    req = requests.get('https://stackoverflow.com')
    html = req.text
    soup = bs(html, 'html.parser')
    my_titles = soup.select('h3 > a')
    data = []
    for title in my_titles:
        data.append(title.get('href'))
    time.sleep(2)

class abc:
    def all(self):
        complex_task()
        return "This is abc \n what a abc!\n"

class bcd:
    def all(self):
        complex_task()
        return "This is bcd \n what a bcd!\n"

class cde:
    def all(self):
        complex_task()
        return "This is cde \n what a cde!\n"

class ijk:
    def all(self):
        complex_task()
        return "This is ijk \n what a ijk!\n"

def main():
    site_list = [abc(), bcd(), cde(), ijk()]
    
    # 单进程版本(保留原代码)
    # start_time = time.time()
    # messages = ''
    # for site in site_list:
    #     messages += site.all()
    # print(messages)
    # print(f"单进程耗时: {time.time() - start_time}")
    
    # 多进程版本:用Pool实现并行
    start_time = time.time()
    # 创建进程池,默认大小为CPU核心数,也可以手动指定(比如processes=4)
    with mp.Pool() as pool:
        # 用map方法并行执行每个site的all方法
        results = pool.map(lambda site: site.all(), site_list)
    # 合并结果
    messages = ''.join(results)
    print(messages)
    print(f"多进程耗时: {time.time() - start_time}")

if __name__ == "__main__":
    main()

为什么这样修改有效?

  • mp.Pool会自动创建多个进程(默认等于你的CPU核心数),每个进程会独立处理site_list中的一个任务,真正实现并行执行。
  • 你的complex_task里包含网络请求(IO密集型操作)和time.sleep(2),这类操作非常适合用多进程并行处理,能大幅减少总耗时——理论上总耗时会接近单个任务的耗时(约网络请求时间+2秒),而不是单进程的4倍耗时。

可选方案:保留队列架构并启动多消费进程

如果你想继续使用队列模式,需要启动多个message_collector进程,比如:

def main():
    site_list = [abc(), bcd(), cde(), ijk()]
    ps_queue = mp.JoinableQueue()
    messages_dict = mp.Manager().dict()
    messages_dict['value'] = ''
    
    # 启动多个消费进程(比如4个,和任务数一致)
    consumer_count = 4
    for _ in range(consumer_count):
        proc = mp.Process(target=message_collector, args=(ps_queue, messages_dict))
        proc.daemon = True
        proc.start()
    
    start_time = time.time()
    crawler(site_list, ps_queue)
    ps_queue.join()
    print(messages_dict['value'])
    print(f"多进程耗时: {time.time() - start_time}")

# 注意修改message_collector函数,把messages_dict作为参数传入
def message_collector(ps_queue, messages_dict):
    while True:
        site = ps_queue.get()
        messages_dict['value'] += site.all()
        ps_queue.task_done()

这种方式也能实现并行,但代码复杂度比用Pool高很多,所以更推荐第一种方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:44:26