如何将内存缓冲区中的文件从Python发送至Julia
跨进程零磁盘IO传输Pandas DataFrame到Julia(Arrow+TCP实现)
前置依赖安装
- Python端:执行
pip install pandas pyarrow安装所需库,pyarrow负责Arrow格式的内存序列化 - Julia端:在Julia REPL执行
] add Arrow DataFrames安装所需库,Sockets为Julia标准库无需额外安装
1. Python端:将Pandas DataFrame写入内存缓冲区(无磁盘IO)
核心逻辑:用内存字节流对象替代磁盘文件路径,直接将DataFrame序列化为Arrow/Feather格式的二进制字节,全程不写入磁盘。
import pandas as pd import pyarrow.feather as feather import io # 此处替换为你自己的大型Pandas DataFrame # df = pd.read_parquet("your_source_data.parquet") # 初始化内存字节缓冲区 mem_buf = io.BytesIO() # 直接将DataFrame以Feather v2(Arrow IPC)格式写入内存 feather.write_feather(df, mem_buf) # 提取内存中序列化完成的二进制内容,用于后续网络传输 arrow_binary = mem_buf.getvalue() mem_buf.close()
2. Python端:通过TCP发送内存中的Arrow二进制数据
TCP为流协议,必须先发送固定长度的长度标记,避免接收端出现粘包、半包问题:先发送8字节无符号整数标记Arrow数据总长度,再发送完整二进制内容即可。
import socket import struct # 按需修改监听地址和端口,同机传输用127.0.0.1即可 BIND_HOST = "127.0.0.1" BIND_PORT = 12345 with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as sock: sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) sock.bind((BIND_HOST, BIND_PORT)) sock.listen(1) print(f"TCP服务已启动,等待Julia端连接 {BIND_HOST}:{BIND_PORT}") conn, client_addr = sock.accept() with conn: print(f"接收连接来自 {client_addr}") # 发送8字节大端序长度头,标记后续Arrow数据总长度 conn.sendall(struct.pack("!Q", len(arrow_binary))) # 发送完整Arrow二进制数据 conn.sendall(arrow_binary) print("数据发送完成")
3. Julia端:从TCP读取数据并转换为Julia DataFrame
Julia端先读取长度头确定数据总大小,收齐所有字节后直接调用Arrow.jl的内存读取接口解析,无需落地磁盘,直接转换为DataFrames.jl对象。
using Sockets, Arrow, DataFrames # 与Python端配置保持一致 CONNECT_HOST = ip"127.0.0.1" CONNECT_PORT = 12345 # 建立TCP连接 conn = connect(CONNECT_HOST, CONNECT_PORT) # 读取8字节长度头,解析得到Arrow数据总长度 header_bytes = read(conn, 8) total_len = ntoh(reinterpret(UInt64, header_bytes)[1]) # 循环读取直到收满所有数据,处理TCP半包问题 recv_buffer = Vector{UInt8}(undef, total_len) read_total = 0 while read_total < total_len read_total += readbytes!(conn, view(recv_buffer, read_total+1:total_len), total_len - read_total) end close(conn) # 直接从内存字节解析Arrow表,转换为Julia DataFrame julia_df = DataFrame(Arrow.Table(recv_buffer)) # 后续可直接使用julia_df对象,类型与Python端自动对齐 println("数据接收完成,共 $(nrow(julia_df)) 行,$(ncol(julia_df)) 列")
性能说明
- 全流程无任何磁盘读写操作,序列化、网络传输、反序列化均为内存操作,GB级DataFrame传输延迟较磁盘临时文件方案低90%以上
- Arrow为跨语言列式格式,类型自动对齐:Pandas的数值、字符串、时间、布尔类型会直接映射为Julia原生类型,无需额外类型转换
- 跨机器传输时仅需将绑定IP改为对应网卡的可访问IP即可,本地回环传输速度接近内存拷贝性能
内容的提问来源于stack exchange,提问作者Jack N
相关产品推荐
相关产品推荐

