如何结合无限循环并发使用ThreadPoolExecutor?多区域数据抓取问题
问题描述
我正在开发一款从相互关联的不同服务器抓取数据的应用。每个服务器对应一个主区域(如美洲、亚洲)及下属子区域(如北美、巴西),且各服务器的子区域相互独立。
由于主区域与子区域的处理方法互不共享,我希望将它们并行运行。我的初步方案是通过executor.map()为主区域创建线程,再在每个主区域线程中为子区域创建线程。所有线程需无限期独立运行。
我的代码如下:
def main_region_fetching(region_list): cur_region = next(iter(region_list.keys())) with concurrent.futures.ThreadPoolExecutor( max_workers=len(region_list["cur_region"]) ) as executor: for region in region_list[cur_region]: executor.submit(sub_region_fetching, region) # 主区域代码应在此处继续执行,但并未到达这里。 with concurrent.futures.ThreadPoolExecutor(max_workers=len(regions)) as executor: executor.map(main_region_fetching, regions)
目前sub_region_fetching会阻塞主线程的执行,请问:
- 它是否会等待线程结束?
- 能否让线程非阻塞运行以实现无限期执行?
- 是否有更优方案,比如使用async?
问题解答
1. 代码是否会等待线程结束?
会的。你用with语句包裹ThreadPoolExecutor时,当with块结束,Executor会自动调用shutdown(wait=True),这会阻塞当前线程直到所有提交的任务执行完毕。如果sub_region_fetching是无限期运行的,这个with块永远不会退出,后面的代码自然无法执行。
另外,外层的executor.map()也会等待所有主区域线程执行完成,才会继续后续逻辑。
2. 如何实现线程非阻塞无限期执行?
要让线程无限期非阻塞运行,核心是避免Executor自动等待任务结束的逻辑:
- 去掉
with语句:直接创建ThreadPoolExecutor实例,不使用上下文管理器。这样Executor不会自动关闭,提交的任务会持续运行。注意程序退出前要手动调用executor.shutdown(wait=False)清理资源,避免泄漏。 - 手动设置
shutdown(wait=False):如果一定要用with,可以在块内手动调用该方法,让with块结束时不等待任务完成,但这种方式灵活性不如直接创建实例。
修改后的示例代码:
def main_region_fetching(region_list): cur_region = next(iter(region_list.keys())) # 直接创建Executor,不使用with executor = concurrent.futures.ThreadPoolExecutor( max_workers=len(region_list[cur_region]) ) for region in region_list[cur_region]: executor.submit(sub_region_fetching, region) # 这里可以继续执行主区域的其他代码,不会被阻塞 # 外层同样不使用with,避免等待所有主区域线程结束 executor = concurrent.futures.ThreadPoolExecutor(max_workers=len(regions)) executor.map(main_region_fetching, regions) # 主线程可继续执行其他逻辑,或保持运行状态
3. 是否有更优方案(比如使用async)?
如果你的抓取任务是IO密集型(比如等待网络响应、服务器返回),使用asyncio异步IO会比多线程更高效——线程在等待IO时会被阻塞,而异步任务在等待时会让出CPU给其他任务,省去线程切换的开销。
异步方案的代码框架示例:
import asyncio async def sub_region_fetching(region): while True: # 执行异步抓取操作(比如用aiohttp发送请求) await fetch_data_from_server(region) # 添加间隔避免频繁请求 await asyncio.sleep(60) async def main_region_fetching(region_list): cur_region = next(iter(region_list.keys())) tasks = [] for region in region_list[cur_region]: task = asyncio.create_task(sub_region_fetching(region)) tasks.append(task) # 主区域可在此执行其他异步操作,或等待子任务(按需选择) await asyncio.gather(*tasks) async def main(): tasks = [] for region_list in regions: task = asyncio.create_task(main_region_fetching(region_list)) tasks.append(task) await asyncio.gather(*tasks) if __name__ == "__main__": asyncio.run(main())
如果任务包含CPU密集型操作,异步IO的优势不明显,这时多线程或多进程会更合适。
内容的提问来源于stack exchange,提问作者Gustavo Feijó
相关产品推荐
相关产品推荐

