ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

PyArrow 与 Java 双向集成实战指南:基于 JPype、pyarrow.jvm 与 C Data Interface 的零拷贝数据交换

PyArrow 与 Java 双向集成实战指南:基于 JPype、pyarrow.jvm 与 C Data Interface 的零拷贝数据交换 PyArrow 与 Java 双向集成实战指南基于 JPype、pyarrow.jvm 与 C Data Interface 的零拷贝数据交换【免费下载链接】arrowApache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing项目地址: https://gitcode.com/gh_mirrors/arrow13/arrow导读Apache Arrow 的核心价值在于让不同语言实现共享同一套内存中的数据。本文基于 Arrow 官方文档 python_java.rst 整理而成系统讲解如何在同一个进程内实现 PythonPyArrow与 JavaArrow Java之间的数据互通从最基础的 JPype 调用 Java 静态方法到使用pyarrow.jvm模块将 JVM 中的 Arrow Vector 直接转换为 PyArrow Array再到基于 C Data Interface含 C Stream Interface实现真正零序列化、零拷贝的双向数据交换。读完本文你将掌握一套完整、可运行的 Python↔Java 混合编程方案并理解其底层的内存共享原理。前置条件与工程准备环境要求本文假设你已经具备以下环境一个正确安装了pyarrow的 Python 环境一个正确安装了arrowArrow Java库的 Java 环境Arrow Java 必须使用mvn -Parrow-c-data编译以确保 C Data 交换支持被启用。为什么需要arrow-c-dataProfile从 java/pom.xml 的源码可以看到arrow-c-data是一个 Maven Profile它额外引入了c模块profile !-- C data interface depends on building a native library -- idarrow-c-data/id modules modulec/module /modules /profile该模块java/cartifactId 为arrow-c-data见 java/c/pom.xml通过 JNI 承载 C Data Interface 的本地实现因此只有在启用该 Profile 构建后Java 端才具备org.apache.arrow.c.*系列类如ArrowArray、ArrowSchema、ArrowArrayStream、Data这是后续 C Data 交换方案的前提。需要安装的 Python 包包用途安装命令jpype1在 Python 解释器中启动 JVM并提供 Java 类代理pip install jpype1cffiPyArrow 的 C Data 导出/导入功能依赖的 C 绑定库pip install cffi其中jpype1是最基本的桥接库cffi仅在后续使用 C Data Interface 时才是必需的pyarrow.cffi模块底层就是基于 cffi 声明了ArrowSchema、ArrowArray、ArrowArrayStream等 C 结构体见 python/pyarrow/cffi.py。从 Python 调用 Java 方法JPype 入门编写并编译一个简单的 Java 类我们首先验证最基础的 Java↔Python 双向可达性。假设有一个输出整数的简单 Java 类public class Simple { public static int getNumber() { return 4; } }将其保存为Simple.java然后用javac编译为 class 文件$ javac Simple.java编译后得到Simple.class。注意这个类没有任何第三方依赖所以javac直接编译即可无需 Maven。使用 JPype 从 Python 调用JPype 能够在 Python 解释器内部启动一个 Java 虚拟机JVM并将 Java 类以代理对象的形式暴露给 Python。创建一个simple.pyimport jpype from jpype.types import * jpype.startJVM(classpath[./]) Simple JClass(Simple) print(Simple.getNumber())逐行解读jpype.startJVM(classpath[./])启动 JVM并把当前目录加入 classpath使Simple.class可被加载JClass(Simple)根据类名获得 Java 类代理Simple.getNumber()直接调用 Java 静态方法。运行$ python simple.py 4输出4说明 Python 代码已经成功访问了 Java 方法。这是整条集成链路的最小闭环也是后续一切复杂方案的地基。注意本文所有示例都假定 Java 类可以在 classpath 中被找到当类依赖外部 jar 包时需要把这些依赖一并加入classpath下文会专门处理。Java 到 Python使用 pyarrow.jvm 转换 Arrow 数组为什么需要 pyarrow.jvm直接通过 JPype 拿到的 Java 对象只是一个代理对象它不具备 PyArrow 数组的方法与能力如类型系统、计算函数、序列化等。pyarrow.jvm模块的作用就是把这些 JVM 上的 Arrow 对象转换为真正的 PyArrow Python 对象。从源码 python/pyarrow/jvm.py 的模块注释可以明确其设计定位These functions convert the objects holding the metadata, the actual data is not copied at all. This will only work with a JVM running in the same process such as provided through jpype.即只转换元数据metadata数据本身零拷贝并且要求 JVM 与 Python 在同一进程中如 JPype因为py4j这类远程 JVM 方案报告的地址在 Python 进程内不可达。编写依赖 Arrow Java 的类FillTen我们创建一个更复杂的 Java 类FillTen.java它使用 Arrow Java 的BigIntVector构建一个包含 1~10 的数组import org.apache.arrow.memory.RootAllocator; import org.apache.arrow.vector.BigIntVector; public class FillTen { static RootAllocator allocator new RootAllocator(); public static BigIntVector createArray() { BigIntVector intVector new BigIntVector(ints, allocator); intVector.allocateNew(10); intVector.setValueCount(10); FillTen.fillVector(intVector); return intVector; } private static void fillVector(BigIntVector iv) { iv.setSafe(0, 1); iv.setSafe(1, 2); iv.setSafe(2, 3); iv.setSafe(3, 4); iv.setSafe(4, 5); iv.setSafe(5, 6); iv.setSafe(6, 7); iv.setSafe(7, 8); iv.setSafe(8, 9); iv.setSafe(9, 10); } }这段代码展示了 Arrow Java 构建向量的标准套路用RootAllocator分配内存Arrow Java 所有向量都从BufferAllocator派生内存allocateNew(10)预分配 10 个槽位setValueCount(10)声明有效元素个数setSafe(i, value)逐个写入值setSafe会自动检查/扩容是安全写入方法。用 Maven 管理依赖并打包由于FillTen依赖了arrow-vector、arrow-memory、arrow-memory-netty、arrow-c-data等一批 jar 包仅靠javac无法编译。需要创建pom.xml收集依赖project modelVersion4.0.0/modelVersion groupIdorg.apache.arrow.py2java/groupId artifactIdFillTen/artifactId version1/version properties maven.compiler.source8/maven.compiler.source maven.compiler.target8/maven.compiler.target /properties dependencies dependency groupIdorg.apache.arrow/groupId artifactIdarrow-memory/artifactId version8.0.0/version typepom/type /dependency dependency groupIdorg.apache.arrow/groupId artifactIdarrow-memory-netty/artifactId version8.0.0/version typejar/type /dependency dependency groupIdorg.apache.arrow/groupId artifactIdarrow-vector/artifactId version8.0.0/version typepom/type /dependency dependency groupIdorg.apache.arrow/groupId artifactIdarrow-c-data/artifactId version8.0.0/version typejar/type /dependency /dependencies /project要点说明maven.compiler.source/target设为 8保证与 Arrow Java 的字节码级别兼容arrow-memory与arrow-vector以pom类型引入它们本质是聚合 POM实际实现类由arrow-memory-nettyNetty 内存分配实现和arrow-c-data提供版本号请替换为你本地实际使用的 Arrow Java 版本。将FillTen.java放到src/main/java/FillTen.java后执行打包$ mvn package [INFO] Scanning for projects... [INFO] [INFO] ------------------ org.apache.arrow.py2java:FillTen ------------------ [INFO] Building FillTen 1 [INFO] --------------------------------[ jar ]--------------------------------- [INFO] [INFO] --- maven-compiler-plugin:3.1:compile (default-compile) FillTen --- [INFO] Changes detected - recompiling the module! [INFO] Compiling 1 source file to /experiments/java2py/target/classes [INFO] [INFO] --- maven-jar-plugin:2.4:jar (default-jar) FillTen --- [INFO] Building jar: /experiments/java2py/target/FillTen-1.jar [INFO] ------------------------------------------------------------------------ [INFO] BUILD SUCCESS [INFO] ------------------------------------------------------------------------产物为target/FillTen-1.jar。收集所有依赖到单一目录为了让 Python 侧能一次性加载 FillTen 及其全部依赖使用 Maven Dependency Plugin 将依赖复制到dependencies目录$ mvn org.apache.maven.plugins:maven-dependency-plugin:2.7:copy-dependencies -DoutputDirectorydependencies输出类似以下为典型输出实际 jar 版本以你的环境为准[INFO] Copying jsr305-3.0.2.jar to /experiments/java2py/dependencies/jsr305-3.0.2.jar [INFO] Copying netty-common-4.1.72.Final.jar to /experiments/java2py/dependencies/netty-common-4.1.72.Final.jar [INFO] Copying arrow-memory-core-8.0.0-SNAPSHOT.jar to /experiments/java2py/dependencies/arrow-memory-core-8.0.0-SNAPSHOT.jar [INFO] Copying arrow-vector-8.0.0-SNAPSHOT.jar to /experiments/java2py/dependencies/arrow-vector-8.0.0-SNAPSHOT.jar [INFO] Copying arrow-c-data-8.0.0-SNAPSHOT.jar to /experiments/java2py/dependencies/arrow-c-data-8.0.0-SNAPSHOT.jar [INFO] Copying arrow-vector-8.0.0-SNAPSHOT.pom to /experiments/java2py/dependencies/arrow-vector-8.0.0-SNAPSHOT.pom [INFO] Copying jackson-core-2.11.4.jar to /experiments/java2py/dependencies/jackson-core-2.11.4.jar [INFO] Copying jackson-annotations-2.11.4.jar to /experiments/java2py/dependencies/jackson-annotations-2.11.4.jar [INFO] Copying slf4j-api-1.7.25.jar to /experiments/java2py/dependencies/slf4j-api-1.7.25.jar [INFO] Copying arrow-memory-netty-8.0.0-SNAPSHOT.jar to /experiments/java2py/dependencies/arrow-memory-netty-8.0.0-SNAPSHOT.jar [INFO] Copying arrow-format-8.0.0-SNAPSHOT.jar to /experiments/java2py/dependencies/arrow-format-8.0.0-SNAPSHOT.jar [INFO] Copying flatbuffers-java-1.12.0.jar to /experiments/java2py/dependencies/flatbuffers-java-1.12.0.jar [INFO] Copying arrow-memory-8.0.0-SNAPSHOT.pom to /experiments/java2py/dependencies/arrow-memory-8.0.0-SNAPSHOT.pom [INFO] Copying netty-buffer-4.1.72.Final.jar to /experiments/java2py/dependencies/netty-buffer-4.1.72.Final.jar [INFO] Copying jackson-databind-2.11.4.jar to /experiments/java2py/dependencies/jackson-databind-2.11.4.jar [INFO] Copying commons-codec-1.10.jar to /experiments/java2py/dependencies/commons-codec-1.10.jar提示除了手工收集依赖也可以使用maven-assembly-plugin构建一个包含全部依赖的单一 fat jar进一步简化 Python 侧 classpath 配置。在 Python 中转换 Java 数组创建fillten_pyarrowjvm.pyimport jpype import jpype.imports from jpype.types import * # Start a JVM making available all dependencies we collected # and our class from target/FillTen-1.jar jpype.startJVM(classpath[./dependencies/*, ./target/*]) FillTen JClass(FillTen) array FillTen.createArray() print(ARRAY, type(array), array) # Convert the proxied BigIntVector to an actual pyarrow array import pyarrow.jvm pyarray pyarrow.jvm.array(array) print(ARRAY, type(pyarray), pyarray) del pyarray注意 classpath 写法./dependencies/*展开为所有收集的依赖 jar./target/*提供FillTen-1.jar。运行结果ARRAY java class org.apache.arrow.vector.BigIntVector [1, 2, 3, 4, 5, 6, 7, 8, 9, 10] ARRAY class pyarrow.lib.Int64Array [ 1, 2, 3, 4, 5, 6, 7, 8, 9, 10 ]第一行FillTen.createArray()的原始返回值它只是 JVM 对象的代理BigIntVector的代理缺乏 PyArrow Array 的能力与接口第二行经pyarrow.jvm.array(array)转换后得到的真正的pyarrow.lib.Int64Array可以像任何 PyArrow 数组一样使用。del pyarray用于显式释放转换后的数组引用配合底层机制释放对 JVM 内存的持有。pyarrow.jvm 模块的源码级原理理解pyarrow.jvm的关键是知道它只搬运元数据、零拷贝数据。阅读 python/pyarrow/jvm.py 可以梳理出完整的转换链路jvm_buffer零拷贝的核心见 jvm.py#L51-L69。它先通过_JvmBufferNanny对org.apache.arrow.memory.ArrowBuf调用getReferenceManager().retain()增加引用计数防止 Python 使用期间 JVM 侧释放内存然后取出memoryAddress()JVM 内存地址与capacity()容量调用pa.foreign_buffer(address, size, basenanny)构造一个直接引用 JVM 内存的 PyArrow Buffer。Python 对象销毁时nanny 的__del__会执行release()归还引用。类型映射field()jvm.py#L199-L256根据 JVMField的TypeID分发到_from_jvm_int_type、_from_jvm_float_type、_from_jvm_time_type、_from_jvm_timestamp_type、_from_jvm_date_type等私有函数把Int、FloatingPoint、Utf8、Binary、FixedSizeBinary、Bool、Time、Timestamp、Date、Decimal、Null等基础类型一一映射到pa.int8()、pa.string()、pa.timestamp(...)等 PyArrow 类型。array()见 jvm.py#L283-L310。对org.apache.arrow.vector.ValueVector取出其字段类型、缓冲区列表getBuffers(False)、元素个数getValueCount()和空值个数getNullCount()最终以pa.Array.from_buffers(dtype, length, buffers, null_count)组装出 PyArrow 数组——缓冲区的底层内存正是 JVM 的。record_batch()见 jvm.py#L313-L335。从VectorSchemaRoot取出 Schema再对每个列向量调用array()最后用pa.RecordBatch.from_arrays组装成 RecordBatch。当前限制文档与源码双重确认pyarrow.jvm目前能力有限——嵌套类型如 struct不受支持复杂类型Struct、List、FixedSizeList、Union、Dictionary在源码中直接抛NotImplementedError见 jvm.py#L243-L247且仅适用于 JPype 这类同一进程内的 JVM。测试用例 python/pyarrow/tests/test_jvm.py 进一步印证了上述行为test_jvm_buffer验证了转换后 Java buffer 引用计数 1、PyArrow 销毁后引用计数还原生命周期绑定test_jvm_types用参数化方式覆盖了从null到decimal128的全部已支持类型映射test_jvm_array验证了BitVector、IntVector、BigIntVector、Float8Vector、各精度TimeStamp*Vector、DateDayVector等向量与 PyArrow 数组的equals断言test_jvm_string_array被标记为xfail因为from_buffers目前只支持基本类型数组——这也解释了字符串数组等复杂结构暂不可用的原因。Java 与 Python 双向通信使用 C Data Interface 零拷贝交换什么是 C Data InterfaceC Data Interface 是 Arrow 为跨语言、跨运行时交换数据而设计的协议核心价值是无需序列化marshaling、无需复制数据Python 侧通过pyarrow.cffi将Array与Schema以ArrowArray*、ArrowSchema*指针导出Java 侧通过org.apache.arrow.carrow-c-data模块将指针重新包裹成ArrowArray/ArrowSchema对象并用Data.importVector导入为FieldVector双方看到的是同一块物理内存一方修改内容另一方直接可见。注官方文档同时指出未来pyarrow.jvm也会改造成基于 C Data Interface 实现目前它仍是专门为 JPype 编写的独立实现。Python 侧导出数组与 Schema使用 C Data Interface 前需要显式安装cffi$ pip install cffi下面是一个 Python 脚本fillten.py它创建一个 10 元素的 PyArrow 数组全 0通过 C Data Interface 导出后交给 Java 方法填充 1~10再观察原数组内容被 Java 修改import jpype import jpype.imports from jpype.types import * # Init the JVM and make FillTen class available to Python. jpype.startJVM(classpath[./dependencies/*, ./target/*]) FillTen JClass(FillTen) # Create a Python array of 10 elements import pyarrow as pa array pa.array([0]*10) from pyarrow.cffi import ffi as arrow_c # Export the Python array through C Data c_array arrow_c.new(struct ArrowArray*) c_array_ptr int(arrow_c.cast(uintptr_t, c_array)) array._export_to_c(c_array_ptr) # Export the Schema of the Array through C Data c_schema arrow_c.new(struct ArrowSchema*) c_schema_ptr int(arrow_c.cast(uintptr_t, c_schema)) array.type._export_to_c(c_schema_ptr) # Send Array and its Schema to the Java function # that will populate the array with numbers from 1 to 10 FillTen.fillCArray(c_array_ptr, c_schema_ptr) # See how the content of our Python array was changed from Java # while it remained of the Python type. print(ARRAY, type(array), array)代码要点arrow_c.new(struct ArrowArray*)在 C 堆上分配结构体arrow_c.cast(uintptr_t, ...)取得指针整数值array._export_to_c(c_array_ptr)把 PyArrow 数组导出为ArrowArray结构。在 PyArrow 实现中见 python/pyarrow/array.pxi#L1705-L1730该操作最终调用 C 层的ExportArray将数组的缓冲区指针、长度、空值数等写入 C 结构array.type._export_to_c(c_schema_ptr)导出类型Schema把两个指针以long传给 Java 静态方法fillCArray。安全提醒官方文档明确指出直接修改数组内容不是安全操作这里仅为演示目的而做。它之所以恰好能工作是因为数组的大小、类型与空值分布都没有改变只是替换了各槽位的值。Java 侧从 C Data 指针导入 FieldVector原有的fillVector是私有方法且只接受BigIntVector。现在需要新增fillCArray把 C Data 交换来的指针重新变成 Arrow Java 能处理的FieldVectorimport org.apache.arrow.c.ArrowArray; import org.apache.arrow.c.ArrowSchema; import org.apache.arrow.c.Data; import org.apache.arrow.memory.RootAllocator; import org.apache.arrow.vector.FieldVector; import org.apache.arrow.vector.BigIntVector; public class FillTen { static RootAllocator allocator new RootAllocator(); public static void fillCArray(long c_array_ptr, long c_schema_ptr) { ArrowArray arrow_array ArrowArray.wrap(c_array_ptr); ArrowSchema arrow_schema ArrowSchema.wrap(c_schema_ptr); FieldVector v Data.importVector(allocator, arrow_array, arrow_schema, null); FillTen.fillVector((BigIntVector)v); } private static void fillVector(BigIntVector iv) { iv.setSafe(0, 1); iv.setSafe(1, 2); iv.setSafe(2, 3); iv.setSafe(3, 4); iv.setSafe(4, 5); iv.setSafe(5, 6); iv.setSafe(6, 7); iv.setSafe(7, 8); iv.setSafe(8, 9); iv.setSafe(9, 10); } }关键步骤ArrowArray.wrap(c_array_ptr)/ArrowSchema.wrap(c_schema_ptr)把裸指针包装为 Java 对象不复制数据Data.importVector(allocator, arrow_array, arrow_schema, null)导入向量。查看 java/c/src/main/java/org/apache/arrow/c/Data.java#L312-L321 可知importVector会先用importField从ArrowSchema还原出Field再field.createVector(allocator)创建向量最后把ArrowArray的内容装载进去第四个参数是字典 Provider本例传null强转为BigIntVector后复用fillVector写入 1~10。由于导入的向量底层缓冲区与 Python 数组共享同一物理内存写入立即反映到 Python 侧。重新编译并运行$ mvn package $ mvn org.apache.maven.plugins:maven-dependency-plugin:2.7:copy-dependencies -DoutputDirectorydependencies $ python fillten.py ARRAY class pyarrow.lib.Int64Array [ 1, 2, 3, 4, 5, 6, 7, 8, 9, 10 ]注意输出中的类型仍是pyarrow.lib.Int64Array——Java 没有重建一个对象而是直接改写 Python 数组背后的内存。进阶使用 C Stream Interface 交换 RecordBatchReaderC Data Interface 只交换单个 Array/Schema若要交换一批数据流则使用 C Stream Interface它基于ArrowArrayStream结构。下面演示 Java 与 Python 之间双向传递pyarrow.RecordBatchReader。Java 侧导出/导入流PythonInteropDemo.java提供了两个能力读取 Arrow IPC 文件并导出为流、导入流并写入 JSON 文件import java.io.File; import java.nio.file.Files; import java.nio.file.Paths; import org.apache.arrow.c.ArrowArrayStream; import org.apache.arrow.c.Data; import org.apache.arrow.memory.BufferAllocator; import org.apache.arrow.memory.RootAllocator; import org.apache.arrow.vector.ipc.ArrowFileReader; import org.apache.arrow.vector.ipc.ArrowReader; import org.apache.arrow.vector.ipc.JsonFileWriter; public class PythonInteropDemo implements AutoCloseable { private final BufferAllocator allocator; public PythonInteropDemo() { this.allocator new RootAllocator(); } public void exportStream(String path, long cStreamPointer) throws Exception { try (final ArrowArrayStream stream ArrowArrayStream.wrap(cStreamPointer)) { ArrowFileReader reader new ArrowFileReader(Files.newByteChannel(Paths.get(path)), allocator); Data.exportArrayStream(allocator, reader, stream); } } public void importStream(String path, long cStreamPointer) throws Exception { try (final ArrowArrayStream stream ArrowArrayStream.wrap(cStreamPointer); final ArrowReader input Data.importArrayStream(allocator, stream); JsonFileWriter writer new JsonFileWriter(new File(path))) { writer.start(input.getVectorSchemaRoot().getSchema(), input); while (input.loadNextBatch()) { writer.write(input.getVectorSchemaRoot()); } } } Override public void close() throws Exception { allocator.close(); } }要点ArrowArrayStream.wrap(cStreamPointer)包装来自 Python 的流指针Data.exportArrayStream(allocator, reader, stream)将 Java 的ArrowFileReader导出为 C 流对应 Data.java#L229 附近的实现Data.importArrayStream(allocator, stream)将 C 流导入为ArrowReader见 Data.java#L408随后逐批loadNextBatch()写入JsonFileWriter类实现AutoCloseable在close()中释放分配器。Python 侧与 Java 互传 RecordBatchReaderdemo.py使用 JPype 启动 JVM创建 Python 的RecordBatchReader导出为ArrowArrayStream指针交给 Java 写 JSON再把 IPC 文件交给 Java 读回流用_import_from_c还原import tempfile import jpype import jpype.imports from jpype.types import * # Init the JVM and make demo class available to Python. jpype.startJVM(classpath[./dependencies/*, ./target/*]) PythonInteropDemo JClass(PythonInteropDemo) demo PythonInteropDemo() # Create a Python record batch reader import pyarrow as pa schema pa.schema([ (ints, pa.int64()), (strs, pa.string()) ]) batches [ pa.record_batch([ [0, 2, 4, 8], [a, b, c, None], ], schemaschema), pa.record_batch([ [None, 32, 64, None], [e, None, None, h], ], schemaschema), ] reader pa.RecordBatchReader.from_batches(schema, batches) from pyarrow.cffi import ffi as arrow_c # Export the Python reader through C Data c_stream arrow_c.new(struct ArrowArrayStream*) c_stream_ptr int(arrow_c.cast(uintptr_t, c_stream)) reader._export_to_c(c_stream_ptr) # Send reader to the Java function that writes a JSON file with tempfile.NamedTemporaryFile() as temp: demo.importStream(temp.name, c_stream_ptr) # Read the JSON file back with open(temp.name) as source: print(JSON file written by Java:) print(source.read()) # Write an Arrow IPC file for Java to read with tempfile.NamedTemporaryFile() as temp: with pa.ipc.new_file(temp.name, schema) as sink: for batch in batches: sink.write_batch(batch) demo.exportStream(temp.name, c_stream_ptr) with pa.RecordBatchReader._import_from_c(c_stream_ptr) as source: print(IPC file read by Java:) print(source.read_all())运行$ mvn package $ mvn org.apache.maven.plugins:maven-dependency-plugin:2.7:copy-dependencies -DoutputDirectorydependencies $ python demo.py JSON file written by Java: {schema:{fields:[{name:ints,nullable:true,type:{name:int,bitWidth:64,isSigned:true},children:[]},{name:strs,nullable:true,type:{name:utf8},children:[]}]},batches:[{count:4,columns:[{name:ints,count:4,VALIDITY:[1,1,1,1],DATA:[0,2,4,8]},{name:strs,count:4,VALIDITY:[1,1,1,0],OFFSET:[0,1,2,3,3],DATA:[a,b,c,]}]},{count:4,columns:[{name:ints,count:4,VALIDITY:[0,1,1,0],DATA:[0,32,64,0]},{name:strs,count:4,VALIDITY:[1,0,0,1],OFFSET:[0,1,1,1,2],DATA:[e,,,h]}]}]} IPC file read by Java: pyarrow.Table ints: int64 strs: string ---- ints: [[0,2,4,8],[null,32,64,null]] strs: [[a,b,c,null],[e,null,null,h]]解读输出JSON 文件Java 端把 Python 传来的RecordBatchReader流逐批写入 JSON。注意其中VALIDITY位图1表示非空、0表示空与OFFSET偏移表字符串列——这正是 Arrow 列式内存布局的直接体现IPC 文件Java 读取 Python 写入的 Arrow IPC 文件后再以 C 流导出回 Pythonpa.RecordBatchReader._import_from_c将其还原为pyarrow.Table该静态方法位于 python/pyarrow/ipc.pxi 的流导出/导入实现中两批数据含null完整保留。三种方案的选型对比方案数据形态拷贝情况适用场景限制JPype 直接调用Java 原生对象无拷贝代理简单方法互调、标量参数返回的 Vector 只是代理无 PyArrow 能力pyarrow.jvmBigIntVector→pyarrow.Int64Array元数据转换、数据零拷贝同进程内JPype把 Java 向量当 PyArrow 数组用仅支持基本类型struct/list 等复杂类型不支持仅限同进程 JVMC Data / C Stream InterfaceArrowArray/ArrowSchema/ArrowArrayStream指针零序列化、零拷贝任意方向批量交换数组/RecordBatch 流需要cffi修改数组内容需自行保证安全三者并非互斥pyarrow.jvm是现成的便捷转换器C Data Interface 是通用底层协议未来pyarrow.jvm也将迁移到 C Data 实现JPype 则是两者共同依赖的 JVM 桥接层。常见问题与注意事项JVM 无法启动/类找不到确认jpype.startJVM的classpath同时包含./dependencies/*与./target/*或 fat jar确认Simple.class/FillTen-1.jar位于对应目录。pyarrow.jvm转换复杂类型报NotImplementedError这是当前已知限制源码 jvm.py#L243-L247 对 Struct、List、Union、Dictionary 直接抛错请改用 C Data Interface 方案或拆分数据为基本类型字段。C Data 修改数组内容恰好能工作的前提数组的大小、类型、空值分布都不能改变。若要安全地由 Java 构造新数组返回给 Python应让 Java 新建FieldVector再通过 C Data 导出而不是就地改写。Arrow Java 必须启用 C Data Profile没有mvn -Parrow-c-data构建的话Java 侧缺少org.apache.arrow.c包C Data 示例无法编译。内存生命周期pyarrow.jvm.jvm_buffer通过retain/release绑定 Java 与 Python 两侧的生命周期见 python/pyarrow/tests/test_jvm.py 的引用计数测试使用 C Data 时导出方负责在消费完成后释放结构体导入方按协议 move 语义接管内容。不同 JVM 进程问题pyarrow.jvm只适用于 JPype 这类同进程JVMpy4j 等远程 JVM 方案拿到的地址在 Python 进程内不可解引用无法使用本方案。总结通过本文你可以看到一条清晰的由浅入深的 Python↔Java 集成路径JPype 打通基础调用——验证两类语言在同一个进程内可以互相调用pyarrow.jvm实现向量级互转——把 JVM 上的 Arrow Vector 变成真正的 PyArrow 数组元数据转换、数据零拷贝C Data Interface 实现内存级共享——Python 导出、Java 导入或反向两侧操作同一块物理内存C Stream Interface 实现批量流交换——RecordBatchReader双向流动配合 IPC 文件、JSON 落盘等真实业务场景。整套方案的价值在于Arrow 的统一内存布局让跨语言调用免去了序列化/反序列化的开销使 Python 与 Java 的混合型数据管道可以在同一进程内以接近原生性能的方式协同工作。需要进一步深入时可继续阅读本仓库中的 pyarrow.jvm 源码、C Data 测试用例、Java C Data 模块 以及 PyArrow 导出/导入实现。【免费下载链接】arrowApache Arrow is a multi-language toolbox for accelerated data interchange and in-memory processing项目地址: https://gitcode.com/gh_mirrors/arrow13/arrow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表