如何让ThreadPool/ProcessPoolExecutor的Worker保持FTP连接状态?
实现带状态的FTP Worker,兼容ThreadPoolExecutor和ProcessPoolExecutor
要实现每个Worker复用FTP连接、避免重复建立连接的需求,需要针对线程池和进程池的特性分别处理——两者的内存隔离机制不同,直接用类对象的方式在进程池里会遇到序列化问题(FTP连接无法跨进程传递)。以下是可行的兼容实现:
核心方案
- 线程池(ThreadPoolExecutor):利用线程局部存储(
threading.local()),让每个线程维护独立的FTP连接,线程内所有下载任务复用该连接。 - 进程池(ProcessPoolExecutor):通过进程初始化函数,在每个进程启动时建立专属的FTP连接,进程内的所有任务复用这个连接。
完整代码实现
import ftplib import threading from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor import os # 确保下载目录存在 os.makedirs("./downloads", exist_ok=True) # -------------------------- 线程池相关实现 -------------------------- # 线程局部存储:每个线程独立保存自己的FTP连接 _thread_local = threading.local() def get_thread_ftp_connection(server, username, password): """获取当前线程的FTP连接,不存在则新建;连接断开时自动重连""" if not hasattr(_thread_local, 'ftp'): ftp = ftplib.FTP(server) ftp.login(username, password) _thread_local.ftp = ftp else: # 发送NOOP命令检查连接是否存活 try: _thread_local.ftp.voidcmd('NOOP') except (ftplib.error_temp, ftplib.error_perm): # 连接失效,重新建立 _thread_local.ftp = ftplib.FTP(server) _thread_local.ftp.login(username, password) return _thread_local.ftp def download_with_thread_pool(filename, server, username, password): """线程池执行的单文件下载任务""" ftp = get_thread_ftp_connection(server, username, password) try: with open(f"./downloads/{filename}", 'wb') as local_file: ftp.retrbinary(f"RETR {filename}", local_file.write) print(f"线程池下载完成:{filename}") except Exception as e: print(f"线程池下载失败 {filename}:{str(e)}") # -------------------------- 进程池相关实现 -------------------------- # 进程全局变量:每个进程独立保存自己的FTP连接 _process_ftp = None def init_process_ftp(server, username, password): """进程初始化时建立FTP连接""" global _process_ftp try: _process_ftp = ftplib.FTP(server) _process_ftp.login(username, password) except Exception as e: print(f"进程初始化FTP连接失败:{str(e)}") raise def download_with_process_pool(filename): """进程池执行的单文件下载任务""" global _process_ftp try: # 检查连接状态 _process_ftp.voidcmd('NOOP') with open(f"./downloads/{filename}", 'wb') as local_file: _process_ftp.retrbinary(f"RETR {filename}", local_file.write) print(f"进程池下载完成:{filename}") except (ftplib.error_temp, ftplib.error_perm): # 连接失效,尝试重连 _process_ftp = ftplib.FTP(os.environ['FTP_SERVER']) _process_ftp.login(os.environ['FTP_USER'], os.environ['FTP_PASS']) download_with_process_pool(filename) except Exception as e: print(f"进程池下载失败 {filename}:{str(e)}") # -------------------------- 统一调用入口 -------------------------- def bulk_ftp_download(filenames, server, username, password, use_process_pool=False, max_workers=12): """ 批量下载FTP文件 :param filenames: 待下载的远程文件名列表 :param server: FTP服务器地址 :param username: FTP用户名 :param password: FTP密码 :param use_process_pool: 是否使用进程池(默认线程池) :param max_workers: 最大并行数 """ if use_process_pool: # 进程池模式:传入初始化函数和参数,每个进程启动时建立连接 with ProcessPoolExecutor( max_workers=max_workers, initializer=init_process_ftp, initargs=(server, username, password) ) as executor: executor.map(download_with_process_pool, filenames) else: # 线程池模式:每个线程维护独立连接 with ThreadPoolExecutor(max_workers=max_workers) as executor: # 用lambda传递额外参数,也可以用functools.partial executor.map( lambda fname: download_with_thread_pool(fname, server, username, password), filenames ) # -------------------------- 使用示例 -------------------------- if __name__ == "__main__": # 配置信息 FTP_SERVER = "your-ftp-server.com" FTP_USER = "your-username" FTP_PASS = "your-password" FILE_LIST = ["file1.txt", "file2.jpg", "file3.zip"] # 替换为你的文件列表 # 使用线程池下载 bulk_ftp_download(FILE_LIST, FTP_SERVER, FTP_USER, FTP_PASS, use_process_pool=False) # 使用进程池下载 # bulk_ftp_download(FILE_LIST, FTP_SERVER, FTP_USER, FTP_PASS, use_process_pool=True)
关键细节说明
- 线程安全与隔离:
- 线程池用
threading.local()保证每个线程的FTP连接独立,避免ftplib非线程安全带来的问题。 - 进程池因内存完全隔离,每个进程的连接互不干扰,无需担心线程安全问题。
- 线程池用
- 连接存活检测:
- 加入了
NOOP命令检查连接状态,连接断开时自动重连,提升稳定性。
- 加入了
- 兼容性:
- 线程池模式直接传递连接参数,进程池模式通过
initializer初始化连接,完美兼容两种Executor。
- 线程池模式直接传递连接参数,进程池模式通过
- 错误处理:
- 捕获了常见的FTP错误和IO错误,避免单个任务失败导致整个批量下载中断。
内容的提问来源于stack exchange,提问作者Bastiaan
相关产品推荐
相关产品推荐

