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

Ruby中如何流式写入Parquet文件?Rails归档工具开发需求

基于red-parquet/red-arrow实现Rails数据分批流式写入Parquet

你提到的Parquet::ArrowFileWriter是Apache Arrow C++层的内部测试类,red-parquet并未将其暴露为对外可用的API。最新版本的red-parquet提供了Parquet::Writer类来实现流式分批写入,具体实现步骤如下:

1. 初始化Writer(基于首个数据批次的Schema)

先从数据库查询第一批数据,将其转换为Arrow的RecordBatch,以此确定Parquet文件的Schema并初始化Writer:

require "red-arrow"
require "red-parquet"

# 从Rails模型获取第一批归档数据(按时间范围+分页)
first_batch = YourArchiveModel.where(created_at: start_time..end_time).limit(1000).to_a

# 定义Arrow Schema,需与模型字段对应
fields = [
  Arrow::Field.new("id", :int64),
  Arrow::Field.new("content", :string),
  Arrow::Field.new("created_at", :timestamp)
]
schema = Arrow::Schema.new(fields)

# 将ActiveRecord对象转换为Arrow列数组
columns = first_batch.map do |record|
  [record.id, record.content, record.created_at]
end.transpose.map { |col| Arrow::Array.new(col, schema.field(col.index).data_type) }
first_record_batch = Arrow::RecordBatch.new(schema, columns)

# 初始化Parquet Writer,指定输出文件路径和Schema
writer = Parquet::Writer.new("archive_#{start_time.strftime("%Y%m%d")}.parquet", first_record_batch.schema)
writer.write(first_record_batch)

2. 分批写入剩余数据

循环查询剩余批次的数据,转换为同Schema的RecordBatch后写入Writer:

# 用主键范围查询替代offset,提升大数量分批的性能
last_id = first_batch.last.id
loop do
  next_batch = YourArchiveModel.where(created_at: start_time..end_time).where("id > ?", last_id).limit(1000).to_a
  break if next_batch.empty?

  # 转换为RecordBatch
  columns = next_batch.map do |record|
    [record.id, record.content, record.created_at]
  end.transpose.map { |col| Arrow::Array.new(col, schema.field(col.index).data_type) }
  record_batch = Arrow::RecordBatch.new(schema, columns)

  writer.write(record_batch)
  last_id = next_batch.last.id
end

3. 关闭Writer完成归档

所有数据写入完成后,必须关闭Writer以确保Parquet文件的元数据正确写入:

writer.close

简化数据转换的技巧

如果不想手动映射字段,可以使用arrow-activerecord gem(需兼容你的Rails版本),它能直接将ActiveRecord Relation转换为Arrow的RecordBatch:

# 安装gem后直接调用
first_record_batch = YourArchiveModel.where(created_at: start_time..end_time).limit(1000).to_arrow_record_batch

关键注意事项

  • 所有批次的RecordBatch必须与初始化Writer时的Schema完全一致,否则会触发写入异常。
  • 避免使用offset进行分批查询,大offset会导致数据库性能急剧下降,优先使用主键或时间戳的范围查询。
  • 写入过程中需保证程序稳定运行,若中途中断,Parquet文件会损坏,可添加异常捕获和断点续写逻辑。

内容的提问来源于stack exchange,提问作者Andy F

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 16:12:50