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

多线程下Paramiko Transport线程阻塞 求超时后关闭连接方案

问题描述

在Windows Server 2016环境下,使用Python Paramiko通过守护线程(setDaemon=True)+ join超时机制批量SSH连接500余台机器。运行后发现,成功连接并关闭50-70台机器的连接后,部分Paramiko Transport连接出现阻塞,状态显示:

<paramiko.Transport at 0x34d41130 (cipher aes128-ctr, 128 bits) (active; 0 open channel(s))>
当阻塞线程数超过500时,系统抛出can't start new threads错误。需要解决:如何在线程超时后强制关闭这些阻塞的Paramiko Transport连接?

相关代码
def run_sub_list_thread(self):
        try:
            threads = []
            for machin_ip in node_sub_list:
               thread = Thread(target=self.foo, args=[machin_ip])
               thread.setDaemon(True)
               threads.append(thread)
               thread.start()
               timeout = 60 * 20 / len(threads)
            for thread in threads:
                thread.join(timeout=timeout)
        except Exception as e:
            print(repr(e))
        return 'yes', 'Success'


def _connect(self, machin_ip):
    _ssh_client = paramiko.SSHClient()
    _ssh_client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
    try:
        _ssh_client.connect(
            machin_ip, port=22, username=password[0], password=''
        )
        _invoked_shell = (
            _ssh_client.invoke_shell()
        )
        self._invoked_shell = _invoked_shell
        t.sleep(2)
    except paramiko.ssh_exception.SSHException as e:
        exception_message = "Exception::{0}::{1}".format(method_name, e)
        print(exception_message)
        return "error", exception_message, None, None
    except Exception as e:
        exception_message = "Exception::{0}::{1}".format(method_name, e)
        print(exception_message)
        return "error", exception_message, None, None
    #Some code here to connect machine properly 
        self._ssh_client = _ssh_client
        self._working_username = _working_username
        self._working_password = _working_password
        return _status, _invoked_shell, _working_username, _working_password
    except Exception as e:
        exception_message = "Exception::{0}::{1}".format(method_name, e)
        print(exception_message)
        return "error", exception_message, None, None

def _disconnect(self):
    method_name = inspect.stack()[0][3]
    class_name = self.__class__.__name__
    try:
        self._invoked_shell.close()
        self._ssh_client.close()
        return 'yes', 'Connection Closed'
    except Exception as e:
        exception_message = f"Exception::{class_name}::{method_name}::{repr(e)}"
        print(exception_message)
        return 'error', exception_message

def foo(self,machin_ip):
    method_name = inspect.stack()[0][3]
    class_name = self.__class__.__name__
    try:
        self._connect(machin_ip)
        # sleep instead of run commands and processing commands output
        t.sleep(300)
        self._disconnect()
        return 'yes', 'Done'
    except Exception as e:
        exception_message = f"Exception::{class_name}::{method_name}::{repr(e)}"
        print(exception_message)
        return 'error', exception_message
解决方案

1. 避免共享实例变量,为每个线程独立维护SSH资源

当前代码中self._ssh_client、self._invoked_shell是实例级变量,多线程共享会导致资源混乱,部分线程无法正确获取自身的连接对象。修改方式:

  • 在foo方法内创建独立的SSH客户端和shell对象,不使用实例变量传递
  • 将_connect改为返回SSH客户端和shell对象,而非赋值给self

修改后的核心代码示例:

def _connect(self, machin_ip):
    _ssh_client = paramiko.SSHClient()
    _ssh_client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
    try:
        _ssh_client.connect(
            machin_ip, port=22, username=password[0], password=''
        )
        _invoked_shell = _ssh_client.invoke_shell()
        t.sleep(2)
        return "success", _invoked_shell, _ssh_client
    except paramiko.ssh_exception.SSHException as e:
        exception_message = f"Exception::_connect::{e}"
        print(exception_message)
        return "error", exception_message, None
    except Exception as e:
        exception_message = f"Exception::_connect::{e}"
        print(exception_message)
        return "error", exception_message, None

def foo(self, machin_ip):
    try:
        status, shell, ssh_client = self._connect(machin_ip)
        if status != "success":
            return 'error', shell
        # 模拟业务操作
        t.sleep(300)
        # 关闭当前线程的独立资源
        shell.close()
        ssh_client.close()
        return 'yes', 'Done'
    except Exception as e:
        exception_message = f"Exception::foo::{repr(e)}"
        print(exception_message)
        # 异常时强制清理资源
        if shell:
            shell.close()
        if ssh_client:
            ssh_client.close()
        return 'error', exception_message

2. 给Paramiko连接全流程设置超时,从源头避免阻塞

在连接、认证、交互等阶段都设置超时,防止无限阻塞:

# _connect方法中添加超时参数
_ssh_client.connect(
    machin_ip, port=22, username=password[0], password='',
    timeout=10,  # 连接超时10秒
    banner_timeout=10,  # 服务器Banner响应超时
    auth_timeout=10  # 身份认证超时
)
# 对shell交互设置超时
_invoked_shell.settimeout(30)

3. 线程超时后主动发送终止信号并强制清理资源

join(timeout)仅等待线程结束,不会强制终止。通过以下方式处理:

  • 用threading.Event作为线程终止信号,让线程主动检查并退出
  • 主线程记录每个线程的SSH资源引用,对超时未结束的线程,强制关闭其Transport连接

核心实现示例:

def run_sub_list_thread(self):
    try:
        threads = []
        thread_resources = {}
        stop_events = {}
        # 创建线程锁保证资源字典的线程安全
        lock = threading.Lock()

        for machin_ip in node_sub_list:
            stop_event = threading.Event()
            thread = Thread(target=self.foo, args=[machin_ip, stop_event, lock])
            thread.setDaemon(True)
            threads.append(thread)
            stop_events[thread.ident] = stop_event
            thread.start()
        
        # 设置单线程超时时间(业务时间+预留缓冲)
        single_thread_timeout = 320
        for thread in threads:
            thread.join(timeout=single_thread_timeout)
            if thread.is_alive():
                # 发送终止信号
                stop_events[thread.ident].set()
                # 强制关闭Transport连接
                with lock:
                    if thread.ident in thread_resources:
                        ssh_client = thread_resources[thread.ident]
                        transport = ssh_client.get_transport()
                        if transport and transport.is_active():
                            transport.close()
    except Exception as e:
        print(repr(e))
    return 'yes', 'Success'

# 修改foo方法,支持终止信号和资源记录
def foo(self, machin_ip, stop_event, lock):
    ssh_client = None
    shell = None
    try:
        status, shell, ssh_client = self._connect(machin_ip)
        if status != "success":
            return 'error', status
        
        # 记录当前线程的SSH资源
        with lock:
            thread_resources[threading.get_ident()] = ssh_client
        
        # 模拟业务操作,定期检查终止信号
        for _ in range(30):
            if stop_event.is_set():
                break
            t.sleep(10)
        
        # 清理资源
        if shell:
            shell.close()
        if ssh_client:
            ssh_client.close()
        return 'yes', 'Done'
    except Exception as e:
        exception_message = f"Exception::foo::{repr(e)}"
        print(exception_message)
        # 异常时强制清理资源
        if shell:
            shell.close()
        if ssh_client:
            ssh_client.close()
        return 'error', exception_message

4. 限制并发线程数,避免系统资源耗尽

一次性启动500个线程会占用大量系统资源,容易触发can't start new threads错误。使用线程池限制并发数:

from concurrent.futures import ThreadPoolExecutor

def run_sub_list_thread(self):
    try:
        # 根据服务器配置调整并发数,比如设置为50
        with ThreadPoolExecutor(max_workers=50) as executor:
            futures = [executor.submit(self.foo, machin_ip) for machin_ip in node_sub_list]
            for future in futures:
                try:
                    # 设置单任务超时
                    result = future.result(timeout=320)
                except Exception as e:
                    print(f"Task failed: {repr(e)}")
                    # 取消未完成的任务
                    future.cancel()
    except Exception as e:
        print(repr(e))
    return 'yes', 'Success'

线程池会自动管理线程的创建与销毁,避免一次性创建过多线程导致系统资源耗尽。


内容的提问来源于stack exchange,提问作者Abdelrahman Gamel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 00:45:56