mpi4py出现‘Message from rank 1 and tag 14 truncated’错误求助
MPI_Recv消息截断错误排查与修复建议
问题场景
在集成mpi4py前编写的测试代码如下,运行时触发消息截断错误:
# mpirun -n 5 python3 MPI_function.py import numpy as np from mpi4py import MPI comm = MPI.COMM_WORLD rank = comm.Get_rank() size = comm.Get_size() np.random.seed(0) N = 21345 # arbitrary values = np.random.rand(N, 1000) if rank == 0: section = np.zeros(int(N / size + N % size)) start = 0 else: # rank != 0: section = np.zeros(int(N / size)) start = rank * len(section) + N % size for i in range(len(section)): section[i] = np.mean(values[start + i]) # print(start + i) if rank != 0: comm.Send(section, dest = 0, tag = 14) else: # rank == 0: results = np.pad(section, (0, N - len(section)), constant_values = 0) for r in range(1, size): temp = np.zeros(N % size) comm.Recv(temp, source = r, tag = 14) start = r * len(N % size) + N % size for i in range(N % size): results[start + i] = temp[i]
运行后报错:
Traceback (most recent call last): File "[FILEPATH]/MPI_function.py", line 33, in <module> comm.Recv(temp, source = r, tag = 14) File "mpi4py/MPI/Comm.pyx", line 299, in mpi4py.MPI.Comm.Recv mpi4py.MPI.Exception: Message truncated, error stack: internal_Recv(127).......: MPI_Recv(buf=0x13be04970, count=0, MPI_DOUBLE, 1, 14, MPI_COMM_WORLD, status=0x1) failed MPIDIG_recv_type_init(77): Message from rank 1 and tag 14 truncated; 0 bytes received but buffer size is 34152 =================================================================================== = BAD TERMINATION OF ONE OF YOUR APPLICATION PROCESSES = PID 31516 RUNNING AT Jacob-Ivanovs-MacBook-Air.local = EXIT CODE: 9 = CLEANING UP REMAINING PROCESSES = YOU CAN IGNORE THE BELOW CLEANUP MESSAGES =================================================================================== YOUR APPLICATION TERMINATED WITH THE EXIT STRING: Killed: 9 (signal 9) This typically refers to a problem with your application. Please see the FAQ page for debugging suggestions
错误核心原因
- 接收缓冲区大小不匹配:非0进程发送的
section长度是N//size(即21345//5=4269),但主进程(rank0)创建的接收缓冲区temp大小是N%size(21345%5=0),缓冲区长度为0,完全无法接收非0进程发来的数据,触发截断错误。 - 结果数组索引计算错误:主进程中计算非0进程结果存放位置的
start公式错误,会导致数组越界或数据覆盖问题。
修复后的代码
# mpirun -n 5 python3 MPI_function.py import numpy as np from mpi4py import MPI comm = MPI.COMM_WORLD rank = comm.Get_rank() size = comm.Get_size() np.random.seed(0) N = 21345 # arbitrary values = np.random.rand(N, 1000) # 计算每个进程处理的数据量:rank0多处理N%size个 chunk_size = N // size if rank == 0: local_size = chunk_size + N % size start_idx = 0 else: local_size = chunk_size start_idx = (N % size) + rank * chunk_size # 计算本地数据的均值 local_means = np.array([np.mean(values[start_idx + i]) for i in range(local_size)]) if rank != 0: # 发送本地计算的均值数组到rank0 comm.Send(local_means, dest=0, tag=14) else: # 初始化结果数组,先放入rank0自己的计算结果 results = np.zeros(N) results[:local_size] = local_means # 接收其他进程的数据并填充到结果数组 for r in range(1, size): # 接收缓冲区大小对应非0进程的chunk_size recv_buf = np.zeros(chunk_size) comm.Recv(recv_buf, source=r, tag=14) # 计算当前进程数据在结果数组中的起始位置 recv_start = (N % size) + r * chunk_size results[recv_start:recv_start + chunk_size] = recv_buf # 验证:rank0打印前10个结果(可选) if rank == 0: print(results[:10])
关键修复点说明
- 缓冲区大小对齐:主进程接收非0进程数据时,缓冲区大小设置为
chunk_size(即N//size),和发送方的local_size一致,避免截断。 - 索引计算修正:非0进程处理数据的起始位置和主进程接收后存放的起始位置统一为
(N%size) + rank*chunk_size,确保数据在结果数组中位置正确。 - 代码可读性优化:变量名更清晰(如
local_size、start_idx),用切片操作替代循环赋值,提升效率和可读性。
额外优化建议
避免手动拆分数据:mpi4py提供Scatterv和Gatherv函数,可以自动完成数据的拆分与聚合,减少手动计算索引的错误,代码更简洁。示例:
# 用Gatherv替代手动Send/Recv的简化版本 if rank == 0: # 定义每个进程发送的数据量和偏移量 sendcounts = np.full(size, chunk_size, dtype=int) sendcounts[0] += N % size displacements = np.concatenate([[0], np.cumsum(sendcounts[:-1])]) results = np.zeros(N) else: sendcounts = None displacements = None results = None comm.Gatherv(sendbuf=local_means, recvbuf=(results, sendcounts, displacements, MPI.DOUBLE), root=0)
内容的提问来源于stack exchange,提问作者Jacob Ivanov
相关产品推荐
相关产品推荐

