如何通过PyArrow在进程间零拷贝共享DataFrame
问题原因
你遇到段错误的核心原因是进程的虚拟地址空间是独立的:第一个进程打印的0x7fa5be7d5850是它自己的虚拟内存地址,在另一个进程中这个地址指向的是完全不同的内存区域(甚至可能未被分配),直接访问必然导致段错误。PyArrow的foreign_buffer只能用于访问当前进程内的外部内存(比如C扩展分配的内存),不支持跨进程裸地址访问。
可行方案
要实现无序列化/磁盘写入的进程间数据共享,需要借助共享内存机制,以下是两种适合你的方案:
方案一:Python标准库multiprocessing.shared_memory + PyArrow
利用Python自带的共享内存模块创建跨进程可见的内存区域,再用PyArrow Buffer包装它,实现零拷贝访问:
进程1(创建共享内存并写入数据)
import time import pyarrow as pa from multiprocessing import shared_memory data = b'abcdefghijklmnopqrstuvwxyz' # 创建共享内存,大小匹配数据长度 shm = shared_memory.SharedMemory(create=True, size=len(data)) # 把数据写入共享内存 shm.buf[:len(data)] = data # 打印共享内存名称,供另一个进程使用 print(f"共享内存名称: {shm.name}") # 保持进程运行,让另一个进程有时间访问 try: while True: time.sleep(1) finally: shm.close() shm.unlink()
进程2(访问共享内存)
import pyarrow as pa from multiprocessing import shared_memory # 替换成进程1打印的共享内存名称 shm_name = "你的共享内存名称" shm = shared_memory.SharedMemory(name=shm_name) # 用PyArrow包装共享内存的buffer buf = pa.py_buffer(shm.buf) print(buf.to_pybytes()) # 输出: b'abcdefghijklmnopqrstuvwxyz' shm.close()
方案二:PyArrow Plasma(适合共享复杂对象如pandas DataFrame)
Plasma是PyArrow官方提供的内存对象存储,专门用于高效的进程间对象共享,支持直接存储Arrow Table(pandas DataFrame可以轻松转为Arrow Table),全程无需序列化到磁盘。
步骤1:启动Plasma存储服务
首先需要启动Plasma后台服务(终端执行):
plasma_store -m 1000000000 -s /tmp/plasma
-m指定可用内存大小(这里是1GB)-s指定Plasma的套接字路径
进程1(存入数据/Arrow Table)
import pyarrow as pa import pyarrow.plasma as plasma import pandas as pd # 连接到Plasma存储 client = plasma.connect("/tmp/plasma") # 示例:创建pandas DataFrame并转为Arrow Table df = pd.DataFrame({"col1": [1,2,3], "col2": ["a","b","c"]}) table = pa.Table.from_pandas(df) # 生成对象ID,存入Plasma object_id = plasma.ObjectID.from_random() client.put(table, object_id) print(f"对象ID: {object_id.binary().hex()}") # 保持连接 try: while True: pass finally: client.disconnect()
进程2(读取共享的对象)
import pyarrow as pa import pyarrow.plasma as plasma import pandas as pd # 连接到Plasma存储 client = plasma.connect("/tmp/plasma") # 替换成进程1打印的对象ID十六进制字符串 object_id_hex = "你的对象ID十六进制字符串" object_id = plasma.ObjectID(bytes.fromhex(object_id_hex)) # 获取Arrow Table并转回pandas DataFrame table = client.get(object_id) df = table.to_pandas() print(df) client.disconnect()
关键说明
- 对于pandas DataFrame,优先用方案二的Plasma,因为它原生支持Arrow格式的对象共享,能最大程度避免数据拷贝。
- 无论哪种方案,核心都是利用操作系统提供的共享内存机制,而不是直接使用进程内的虚拟地址。
内容的提问来源于stack exchange,提问作者orange
相关产品推荐
相关产品推荐

