Python中使用信号量时如何正确实现线程join?
解决vCenter 6.7线程过载问题:限制并发数并正确等待所有任务完成
核心问题分析
你当前的代码一次性启动2000+线程并发请求vCenter 6.7,超出了旧版本vCenter的连接承载上限,导致连接断开、线程崩溃。需要限制同时运行的线程数量,同时确保所有采集任务完成后再进入下一环节。
可行方案
方案1:使用concurrent.futures.ThreadPoolExecutor(推荐)
线程池会自动管理并发线程数,无需手动创建和join线程,还能方便地控制超时和任务结果。
修改后的代码示例:
from concurrent.futures import ThreadPoolExecutor, as_completed try: objects = vsphere_client.vcenter.VM.list() # 尝试新版REST API old_objects = container_view.view # 旧版pyvmomi对象 rest_api = True except UnableToAllocateResource: # vCenter 6.7在1000个对象时触发该错误,7.0为4000个 objects = container_view.view old_objects = None except OperationNotFound: objects = container_view.view old_objects = None # 定义每个VM的采集任务 def collect_vm_detail(obj): collector = RESTVMDetail(vsphere_client, db_vcenter, obj, old_objects, rest_api, db_vms, db_hosts, db_datastores, db_networks, db_vm_disks, db_vm_os_disks, db_vm_nics, db_vm_cdroms, db_vm_floppies, db_vm_scsis, db_regions, db_sites, db_environments, db_platforms, db_applications, db_functions, db_costs, db_vm_snapshots, api_limiter) collector.run() # 直接执行采集逻辑,线程池负责管理线程 # 设置并发数(根据vCenter 6.7的承载能力调整,比如设为200) MAX_WORKERS = 200 with ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor: # 提交所有采集任务 futures = [executor.submit(collect_vm_detail, obj) for obj in objects] # 等待任务完成,逐个处理超时与异常 for future in as_completed(futures): try: future.result(timeout=600) # 单个任务10分钟超时 except Exception as e: print(f"VM采集任务超时或失败: {e}")
方案2:基于原线程类使用threading.Semaphore
如果要保留手动创建线程的方式,用信号量限制同时运行的线程数:
首先修改RESTVMDetail线程类,加入信号量控制:
import threading class RESTVMDetail(threading.Thread): def __init__(self, semaphore, *args, **kwargs): super().__init__(*args, **kwargs) self.semaphore = semaphore # 原有初始化逻辑... def run(self): with self.semaphore: # 自动完成信号量的获取与释放 # 原有VM数据采集逻辑...
然后在主代码中创建信号量并传递给每个线程:
# 设置最大并发数(根据vCenter承载能力调整) MAX_CONCURRENT = 200 semaphore = threading.Semaphore(MAX_CONCURRENT) threads = [] for obj in objects: thread = RESTVMDetail(semaphore, vsphere_client, db_vcenter, obj, old_objects, rest_api, db_vms, db_hosts, db_datastores, db_networks, db_vm_disks, db_vm_os_disks, db_vm_nics, db_vm_cdroms, db_vm_floppies, db_vm_scsis, db_regions, db_sites, db_environments, db_platforms, db_applications, db_functions, db_costs, db_vm_snapshots, api_limiter) threads.append(thread) # 启动所有线程 for thread in threads: thread.start() # 等待所有线程完成,单个线程超时10分钟 for thread in threads: thread.join(600)
方案说明
- 线程池方案更简洁,无需手动管理线程生命周期,
as_completed可灵活处理单个任务的超时和异常,避免整体流程被卡住。 - 信号量方案保留了你原有的线程类结构,通过
with self.semaphore确保同一时间只有指定数量的线程执行采集逻辑,不会一次性压垮vCenter。 - 两种方案都能确保所有任务完成后再进入下一环节,解决了你之前join线程的阻塞困扰——线程池的
shutdown(wait=True)(with语句自动调用)或手动join所有线程,都能等待全部任务结束。
内容的提问来源于stack exchange,提问作者Jimmy Fort
相关产品推荐
相关产品推荐

