Python ThreadPool线程串行执行原因及保留现有写法的修复方案
问题描述
更新记录
- 更新1:如果将for循环内的代码修改为:
print('processing new page') pool.apply_async(time.sleep, (5,))
观察到每一次打印后都会出现5秒延迟,因此该问题与webdriver无关。
- 更新2:感谢@user56700的解答,但希望了解当前写法的问题所在,且在不改变现有线程使用方式的前提下修复问题。
问题背景
最初编写的Python串行代码如下:
driver = webdriver.Chrome(options=chrome_options, service=Service('./chromedriver')) for url in url_list: try: print('processing new page') result = parse_page(driver, url) # 通过driver访问url,等待页面加载完成后解析内容(单页耗时约30秒) # 修改全局变量 except Exception as e: log_warning(str(e))
如果处理10个页面,上述代码总耗时达300秒,效率极低。
尝试引入多线程优化,但实现后未达到预期效果,多线程实现代码如下:
import threading from multiprocessing.pool import ThreadPool as Pool G_LOCK = threading.Lock() driver = webdriver.Chrome(options=chrome_options, service=Service('./chromedriver')) pool = Pool(10) for url in url_list: try: print('processing new page') result = pool.apply_async(parse_page, (driver, url,)).get() G_LOCK.acquire() # 修改全局变量 G_LOCK.release() except Exception as e: log_warning(str(e)) pool.close() pool.join() # 此处需确保所有线程执行完成后再运行后续代码
实现全程复用同一个driver实例,实际运行时在processing new page打印语句旁添加时间戳后,输出如下:
[10:36:02] processing new page [10:36:09] processing new page [10:36:15] processing new page [10:36:22] processing new page [10:36:39] processing new page
该输出不符合预期,预期两次打印间隔应在1秒左右,因为线程提交后仅需修改全局变量即可。
问题根因
当前写法存在两个核心错误:
- 异步任务提交后立刻调用
get()造成阻塞,完全没有实现并发pool.apply_async()的作用是把任务提交到线程池异步执行,调用后会立刻返回异步结果对象,不会阻塞主线程。但你在提交任务后立刻调用了.get()方法,这个方法会阻塞当前主线程,直到刚提交的任务完全执行结束、拿到返回值,才会进入下一轮循环提交下一个任务。这种写法本质还是“提交一个任务→等任务跑完→再提交下一个”,和串行执行逻辑完全一致,这也是你哪怕把任务替换成time.sleep(5),依然会每打印一次就等待5秒的根本原因,和webdriver本身没有关系。 - 复用单个Selenium driver实例本身不支持多线程并发操作
Selenium生成的driver实例不是线程安全的,就算修复了上述阻塞问题,多个线程同时操作同一个driver访问不同URL,会出现页面抢占、会话混乱、元素定位失败等不可预期的错误。
修复方案(不改变现有线程使用方式)
保留ThreadPool的使用逻辑,按以下步骤修复:
- 拆分任务提交和结果获取逻辑:循环中只负责把所有任务提交到线程池,保存所有返回的异步结果对象,不要在提交循环里调用
get()造成阻塞 - 所有任务提交完成后,再统一遍历异步结果对象,等待任务执行完成、获取返回值,再处理全局变量的修改
- 由于需要复用同一个driver实例,必须额外给driver的操作逻辑加独立锁,保证同一时间只有一个线程能操作driver,避免并发冲突。
修复后的参考代码:
import threading from multiprocessing.pool import ThreadPool as Pool G_LOCK = threading.Lock() # 新增driver操作锁,避免多线程同时操作同一个driver引发异常 DRIVER_LOCK = threading.Lock() driver = webdriver.Chrome(options=chrome_options, service=Service('./chromedriver')) pool = Pool(10) async_tasks = [] def wrapped_parse_task(url): # 加锁保证同一时刻只有一个线程操作driver with DRIVER_LOCK: print('processing new page') parse_result = parse_page(driver, url) return parse_result for url in url_list: try: # 仅提交任务,不等待结果,主线程会立刻进入下一轮循环 task = pool.apply_async(wrapped_parse_task, (url,)) async_tasks.append(task) except Exception as e: log_warning(str(e)) # 所有任务提交完成后,统一等待结果并修改全局变量 for task in async_tasks: try: result = task.get() with G_LOCK: # 执行修改全局变量的逻辑 pass except Exception as e: log_warning(str(e)) pool.close() pool.join() # 所有任务执行完成后,再执行后续逻辑
说明:由于要求复用同一个driver实例,上述代码中给driver操作加了全局锁,实际driver访问页面、加载内容的逻辑依然是串行执行的,仅页面解析等非driver操作可以并行,能节省一部分时间。如果要实现完全并发提效,需要在每个工作线程中初始化独立的driver实例,不能全局复用单个driver。
内容的提问来源于stack exchange,提问作者john
相关产品推荐
相关产品推荐

