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

如何通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 03:50:27