如何在不同Python版本与脚本间共享大型DataFrame?
针对你在不同Python(3.8/3.10)、Pandas(1.2/2.0)环境下共享1500万行级DataFrame的需求,以下是无需额外格式转换、低开销的可行方案:
1. Apache Arrow 零拷贝IPC传输(推荐)
Apache Arrow是跨语言、跨版本的列式内存标准,Pandas 1.2+和2.0都原生支持与Arrow格式的互转,且能实现零拷贝数据共享,完美适配双向通信需求。
核心思路
将DataFrame转为Arrow的内存格式,通过Socket流或共享内存在脚本间传输,接收端直接转回DataFrame,全程无需转换为字典或其他中间格式,开销极低。
代码示例(Socket双向传输)
发送端(如Script A向B传数据)
import pandas as pd import pyarrow as pa import socket def send_df_via_socket(df, host="localhost", port=9999): # 将DataFrame转为Arrow表 arrow_table = pa.Table.from_pandas(df) # 创建Socket连接 with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s: s.connect((host, port)) # 序列化Arrow表并发送 writer = pa.ipc.new_stream(s, arrow_table.schema) writer.write_table(arrow_table) writer.close() # 示例:发送ArcticDB中读取的DataFrame # df = arctic_lib.read("data_store") # send_df_via_socket(df)
接收端(如Script B接收A的数据)
import pandas as pd import pyarrow as pa import socket def recv_df_via_socket(host="localhost", port=9999): with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s: s.bind((host, port)) s.listen() conn, addr = s.accept() with conn: # 读取Arrow流并转为DataFrame reader = pa.ipc.open_stream(conn) arrow_table = reader.read_all() df = arrow_table.to_pandas() return df # 示例:接收数据并处理 # df = recv_df_via_socket()
该方式支持双向通信,且完全兼容不同Python/Pandas版本——Arrow的序列化格式不依赖语言或库版本,仅需保证两端都安装pyarrow(10+版本即可兼容Pandas 1.2和2.0)。
2. 基于ArcticDB的Arrow格式存储优化
你已在使用ArcticDB,可将其存储格式配置为Arrow,这样不同版本的Pandas都能直接读取/写入ArcticDB中的数据,无需额外转换:
配置方式
在Script A初始化ArcticDB时指定存储格式为Arrow:
from arctic import Arctic from arctic.store.version_store import VersionStore import pyarrow as pa # 初始化ArcticDB并设置Arrow序列化 ac = Arctic("mongodb://localhost:27017/") ac.initialize_library("data_lib", lib_type=VersionStore) lib = ac["data_lib"] # 设置用Arrow序列化DataFrame lib.set_config("serialization.version", "2") # 启用Arrow序列化 lib.set_config("serialization.format", "arrow")
之后,Script B(Pandas 1.2)和Script C都能直接通过ArcticLib读取DataFrame,ArcticDB会自动处理版本兼容问题,且数据存储/读取的开销远低于手动转换格式。
3. Dask分布式内存共享(适合Dask DataFrame场景)
如果你的数据是Dask DataFrame,可通过Dask分布式集群实现跨脚本的内存共享:
步骤
- 在任意环境启动Dask集群:
dask scheduler
- 各脚本连接同一集群,直接共享Dask DataFrame:
# Script A: 将数据存入集群内存 from dask.distributed import Client import dask.dataframe as dd client = Client("tcp://localhost:8786") # 从ArcticDB读取转为Dask DataFrame # ddf = dd.from_pandas(arctic_lib.read("data_store"), npartitions=10) # 存入集群内存 ddf.persist() client.publish_dataset(data=ddf, name="shared_data") # Script B: 读取集群中的共享数据 from dask.distributed import Client client = Client("tcp://localhost:8786") ddf = client.get_dataset("shared_data") # 执行机器学习预测后,将结果存回ArcticDB或重新发布
Dask的分布式存储同样基于Arrow格式,兼容不同Python/Pandas版本,且支持双向的数据更新与共享。
内容的提问来源于stack exchange,提问作者r-stiller

