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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 13:33:32