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

