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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 18:32:37