重新接入tmux会话时Python多进程脚本触发BrokenPipeError问题
问题描述
我在大学实验室服务器上运行长时间的Python脚本,为避免SSH连接中断导致进程终止,操作流程如下:
- 启动tmux会话:
tmux new -s preprocessing
- 在会话内运行脚本:
python3 path/to/my/scirpt.py
- 退出会话并关闭SSH连接,一段时间后重新连接并接入会话:
tmux attach -t preprocessing
刚接入时能看到脚本的进度条输出:
Processing 1st Step: 7%|███▌ | 61725/921600 [11:11:49<155:59:00, 1.53it/s]
但几秒后脚本抛出BrokenPipeError错误:
multiprocessing.pool.RemoteTraceback: """ Traceback (most recent call last): File "/usr/lib/python3.10/multiprocessing/pool.py", line 125, in worker result = (True, func(*args, **kwds)) File "/usr/lib/python3.10/multiprocessing/pool.py", line 51, in starmapstar return list(itertools.starmap(args[0], args[1])) File "/home/plodir2/network-predictive-analysis/scripts/data_preprocessing_MILC+LAMMPS+UR_lammps_1056.py", line 419, in process_chunk if (application_iteration, router_id) not in list(processed_routers): File "<string>", line 2, in __len__ File "/usr/lib/python3.10/multiprocessing/managers.py", line 817, in _callmethod conn.send((self._id, methodname, args, kwds)) File "/usr/lib/python3.10/multiprocessing/connection.py", line 206, in send self._send_bytes(_ForkingPickler.dumps(obj)) File "/usr/lib/python3.10/multiprocessing/connection.py", line 411, in _send_bytes self._send(header + buf) File "/usr/lib/python3.10/multiprocessing/connection.py", line 368, in _send n = write(self._handle, buf) BrokenPipeError: [Errno 32] Broken pipe """ The above exception was the direct cause of the following exception: Traceback (most recent call last): File "/home/plodir2/network-predictive-analysis/scripts/data_preprocessing_MILC+LAMMPS+UR_lammps_1056.py", line 600, in <module> pool.starmap(process_chunk, [(application_chunk, data_to_add_list, processed_routers, inactive_nodes_lists_processed_iterations) for application_chunk in application_chunks]) File "/usr/lib/python3.10/multiprocessing/pool.py", line 375, in starmap return self._map_async(func, iterable, starmapstar, chunksize).get() File "/usr/lib/python3.10/multiprocessing/pool.py", line 774, in get raise self._value BrokenPipeError: [Errno 32] Broken pipe
关键观察
- 报错仅在重新接入tmux会话后发生,进度条显示的11小时正好是脚本启动到我检查进度的时长。
- 用小数据集测试时,保持SSH连接活跃未出现任何错误。
脚本核心逻辑
脚本将pandas的application_data数据集拆分为多块,通过400个进程并行处理以提升性能,使用multiprocessing.Manager创建进程间共享的数据结构:
... if __name__ == "__main__": num_processes = 400 application_chunks = np.array_split(application_data, num_processes) with Manager() as manager: data_to_add_list = manager.list() processed_routers = manager.list() inactive_nodes_lists_processed_iterations = manager.list() progress_counter = Value('i', 0) lock = Lock() with tqdm(total=len(application_data), desc="Processing 1st Step") as pbar: progress_thread = Thread(target=update_progress_bar, args=(len(application_data), pbar)) progress_thread.start() with Pool(num_processes, initializer=init, initargs=(lock, progress_counter)) as pool: pool.starmap(process_chunk, [(application_chunk, data_to_add_list, processed_routers, inactive_nodes_lists_processed_iterations) for application_chunk in application_chunks]) progress_thread.join() ...
其中processed_routers用于避免多个进程重复处理同一(application_iteration, router_id)元组,读写逻辑如下:
def process_chunk(chunk, data_to_add_list, processed_routers, inactive_nodes_lists_processed_iterations): for application_data_record in chunk.itertuples(): application_iteration = application_data_record.iteration # 整数类型 router_id = int(application_data_record.node // 4) ... add_router_data_ports = False with lock: if (application_iteration, router_id) not in list(processed_routers): # 标记该迭代下的路由器端口数据已处理 processed_routers.append((application_iteration, router_id)) add_router_data_ports = True if add_router_data_ports: # 执行具体处理逻辑 ...
解决方案
1. 优化进程间共享数据结构的性能
当前每次判断(application_iteration, router_id)是否已处理时,都会将共享的processed_routers列表转为本地列表,这会触发跨进程的全量数据同步,当列表规模变大后,极易引发连接超时或断开。
改用共享字典存储已处理元组,判断操作的时间复杂度从O(n)降为O(1),大幅减少跨进程通信量:
# 替换Manager.list()为Manager.dict() processed_routers = manager.dict() # 修改process_chunk中的判断逻辑 with lock: key = (application_iteration, router_id) if key not in processed_routers: processed_routers[key] = True # 用字典键存储元组,值可设为任意标识 add_router_data_ports = True
2. 脱离终端依赖避免tmux重连干扰
tmux会话重连时会重置终端的标准输入输出关联,而tqdm进度条依赖终端文件描述符,这可能导致进程间通信的管道异常。可以通过两种方式解决:
- 将输出重定向到日志文件:运行脚本时脱离终端,直接将输出写入文件
python3 path/to/my/scirpt.py > preprocessing.log 2>&1 &
之后查看进度只需读取日志:
tail -f preprocessing.log
- 禁用
tqdm的终端交互模式:在脚本中强制tqdm输出到标准输出,不依赖终端特性
with tqdm(total=len(application_data), desc="Processing 1st Step", disable=None, leave=True) as pbar: # 保持原进度更新逻辑
3. 调整进程数量至合理范围
400个进程远超多数实验室服务器的CPU核心数(通常为32-64核),会导致严重的进程调度开销,加剧跨进程通信的不稳定性。建议将进程数设置为CPU核心数的1-2倍:
import multiprocessing num_processes = multiprocessing.cpu_count() * 2 # 例如32核服务器设置为64
4. 增强跨进程连接的稳定性
对于Python 3.8及以上版本,可以通过自定义Manager启用心跳机制,避免长时间无通信导致的连接断开:
from multiprocessing.managers import SyncManager # 自定义Manager并注册所需数据类型 class CustomManager(SyncManager): pass CustomManager.register('dict') CustomManager.register('list') # 初始化Manager时设置心跳间隔(单位:秒) with CustomManager(address=('', 50000), authkey=b'secret_key', heartbeat_interval=30) as manager: # 原逻辑:创建共享数据结构、启动进程池等
内容的提问来源于stack exchange,提问作者Pietro Lodi Rizzini
相关产品推荐
相关产品推荐

