多线程队列SSH程序移除sleep后出现命令重复执行问题
问题原因分析
你的线程模型设计存在核心缺陷,导致任务分配逻辑混乱:
- 你创建了主机数 × 命令数的线程总量,每个线程绑定固定命令,但线程仅从队列中获取IP,无法区分当前IP需要执行哪条命令。
- 主线程向队列中重复放入单一IP(每个IP被放入命令数次),队列无法关联IP与对应命令的绑定关系。
- 无
sleep时,主线程快速将所有IP推入队列,大量线程同时争抢队列元素,会出现多个绑定同一命令的线程抢到同一个IP,而绑定其他命令的线程错失该IP,最终引发部分命令重复执行、部分命令缺失的异常。 - 添加
sleep只是通过延迟任务入队,给不同命令的线程留出抢IP的窗口,属于侥幸实现预期效果的临时方案,完全不可靠。
不依赖sleep的解决方法
重构线程模型,让队列存储**(主机IP, 要执行的命令)**完整任务元组,工作线程从队列取出完整任务后执行,从根源上避免任务混淆。
1. 修改任务执行函数
让线程处理任意完整任务,不再绑定固定命令:
def launcher(q): """从队列获取完整任务(IP+命令)并执行""" while True: ip, cmd = q.get() print(f"Thread {ip}: Running {cmd} to {ip}\n") try: subprocess.check_output(f"ssh xxxx@{ip} {cmd}", stderr=subprocess.STDOUT, text=True, shell=True) finally: q.task_done()
2. 重构线程池创建逻辑
创建固定数量的工作线程(最多25个),避免无意义的线程泛滥:
# 创建工作线程池,最多25个线程 num_threads = min(len(ips), 25) for _ in range(num_threads): worker = Thread(target=launcher, args=(queue,)) worker.daemon = True worker.start()
3. 修改任务入队逻辑
将每个主机与命令的组合作为完整任务推入队列:
print("Main Thread Waiting") for ip in ips: for cmd in cmds: queue.put( (ip, cmd) ) # 推入(主机IP, 命令)元组 queue.join()
完整修正代码
#!/usr/bin/env python import subprocess import configparser from threading import Thread from queue import Queue import time """ A threaded ssh based command dispatch system """ start = time.time() queue = Queue() def readConfig(file="config.ini"): """Extract IP addresses and CMDS from config file and returns tuple""" ips = [] cmds = [] Config = configparser.ConfigParser() Config.read(file) machines = Config.items("MACHINES") commands = Config.items("COMMANDS") for ip in machines: ips.append(ip[1]) for cmd in commands: cmds.append(cmd[1]) return ips, cmds def launcher(q): """从队列获取完整任务(IP+命令)并执行""" while True: ip, cmd = q.get() print(f"Thread {ip}: Running {cmd} to {ip}\n") try: subprocess.check_output(f"ssh xxxx@{ip} {cmd}", stderr=subprocess.STDOUT, text=True, shell=True) finally: q.task_done() # 读取配置文件 ips, cmds = readConfig() # 创建工作线程池 num_threads = min(len(ips), 25) for _ in range(num_threads): worker = Thread(target=launcher, args=(queue,)) worker.daemon = True worker.start() # 推入所有任务 print("Main Thread Waiting") for ip in ips: for cmd in cmds: queue.put( (ip, cmd) ) queue.join() end = time.time() print(f"Dispatch Completed in {float(end - start)} seconds")
方案优势
- 每个队列元素都是明确的任务,彻底避免IP与命令的绑定混乱。
- 工作线程数量可控,不会因主机和命令过多创建大量线程浪费资源。
- 依赖队列的线程安全机制分配任务,每个任务仅被执行一次,无需依赖
sleep来维持执行顺序。
内容的提问来源于stack exchange,提问作者parmerxceptn
相关产品推荐
相关产品推荐

