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

Python multiprocessing Pipe大数据发送死锁问题咨询

问题

在Ubuntu系统上使用Python的multiprocessing包开发代码时,为提升进程间通信效率,将待通过multiprocessing.Pipe发送的数据收集到数组中,当元素超过1000个时批量发送。但代码出现偶发死锁,且死锁似乎出现在批量发送数据的环节。

编写测试代码验证后发现:当pickle后的payload超过64KiB时会触发死锁,≤64KiB则可正常传输。但官方文档指出仅当对象大小超过32MiB时才会抛出异常,想咨询这是否为正常情况,或是代码存在问题?

测试代码如下:

import multiprocessing as mp
import sys
import pickle

c, p = mp.Pipe(duplex=False)

batch_list = []

for i in range(21912):
    batch_list.append(i)
    
print(f"Size of payload is: {sys.getsizeof(batch_list)/1024} KiB")
pkl_batch_list = pickle.dumps(batch_list)
print(f"Size of pickled payload is: {sys.getsizeof(pkl_batch_list)/1024} KiB")
p.send(batch_list)


while c.poll():
    try:
        c.recv()
    except Exception:
        print("got exception")

print("Done!")
分析与解决

这不是正常情况,核心问题出在代码逻辑里:Pipe的发送和接收必须在不同进程中执行,而你的测试代码(以及生产代码的逻辑)是在同一个进程里调用p.send()和c.recv(),这会直接导致死锁。

为什么超过64KiB才触发死锁?

Linux系统上multiprocessing.Pipe底层用的是Unix域套接字,当发送的数据≤64KiB时,操作系统套接字缓冲区能一次性容纳全部数据,send()调用会直接返回;但数据超过这个阈值时,缓冲区被填满,send()会进入阻塞状态,等待接收方读取数据腾出空间——但你的代码里发送和接收都在同一个进程,没有其他进程去读取管道,自然就死锁了。

官方文档提到的32MiB限制,是multiprocessing内部对单个pickle对象的大小限制(超过会抛ValueError),和你遇到的死锁是完全不同的问题。

修复方法

必须把发送和接收逻辑放到不同进程中执行,示例代码如下:

import multiprocessing as mp
import sys
import pickle

def sender(pipe):
    batch_list = []
    for i in range(21912):
        batch_list.append(i)
    
    print(f"Size of payload is: {sys.getsizeof(batch_list)/1024} KiB")
    pkl_batch_list = pickle.dumps(batch_list)
    print(f"Size of pickled payload is: {sys.getsizeof(pkl_batch_list)/1024} KiB")
    pipe.send(batch_list)
    pipe.close()

def receiver(pipe):
    while pipe.poll():
        try:
            data = pipe.recv()
            print(f"Received {len(data)} items")
        except Exception as e:
            print(f"Exception occurred: {e}")
    pipe.close()

if __name__ == "__main__":
    c, p = mp.Pipe(duplex=False)
    sender_proc = mp.Process(target=sender, args=(p,))
    receiver_proc = mp.Process(target=receiver, args=(c,))
    
    receiver_proc.start()
    sender_proc.start()
    
    sender_proc.join()
    receiver_proc.join()
    
    print("Done!")

额外注意事项

  • 批量发送数据提升效率的思路是对的,但要确保发送方和接收方在独立进程中运行,避免同一进程内操作管道导致的阻塞。
  • 生产环境中,建议给recv()设置超时时间,避免因异常情况导致接收进程无限阻塞。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 21:13:20