Julia写入Arrow流文件,Python读取仅获首个记录批次问题
问题:Julia写入多批次Arrow流式IPC文件,Python读取仅获首个批次
我通过Julia循环写入1000个包含OutputData结构体的Arrow记录批次至文件,使用Python的PyArrow读取该流式IPC格式文件时,仅能获取首个记录批次,但已确认文件大小符合预期。
原Julia写入代码:
using Arrow struct OutputData name::String age::Int32 end writer = open(filePath, "w") data = OutputData("Alex", 20) for _ = 1:1000 t = (col1=[data],) table = Arrow.Table(Arrow.tobuffer(t)) Arrow.write(writer, table) end close(writer)
Python读取代码:
import pyarrow as pa with pa.OSFile(input_filepath, 'rb') as source: with pa.ipc.open_stream(source) as reader: for batch in reader: # 仅能遍历到一个批次 my_struct_col = batch.column('col1') field1_values = my_struct_col.flatten() print(field1_values)
原因分析与修复
问题根源
原Julia代码每次循环调用Arrow.write(writer, table)时,默认会写入完整的Arrow文件头+批次数据。这导致最终文件是1000个独立的Arrow流拼接而成,而非单一连续的流式IPC序列。PyArrow的ipc.open_stream只会解析第一个完整的流,后续内容被判定为无效数据,因此只能读取到首个批次。
修复后的Julia代码
使用Arrow.Stream创建流式写入器,确保所有批次都追加到同一个IPC流中:
using Arrow struct OutputData name::String age::Int32 end filePath = "output.arrow" writer = open(filePath, "w") stream = Arrow.Stream(writer) # 创建流式写入器 data = OutputData("Alex", 20) for _ = 1:1000 t = (col1=[data],) table = Arrow.Table(t) # 直接从元组创建Table,无需转buffer Arrow.write(stream, table) # 写入到流式writer,共享同一个文件头 end close(stream) close(writer)
关键说明
Arrow.Stream会在首次写入时生成唯一的文件头,后续所有批次都仅写入数据部分,保证文件是单一连续的IPC流。- 移除了不必要的
Arrow.tobuffer调用,直接从元组构造Arrow.Table,简化代码逻辑。
修复后,使用原Python读取代码即可遍历到全部1000个记录批次。
内容的提问来源于stack exchange,提问作者BAR
相关产品推荐
相关产品推荐

