如何用C++的Apache Arrow将动态二维数组写入Parquet文件并解决空文件问题
解决C++动态二维数组写入Parquet文件为空问题并实现指定格式输出
错误原因分析
原代码存在三个关键问题导致生成空Parquet文件:
- 类型不匹配:使用
arrow::FloatBuilder(对应float32类型),但schema定义的字段类型是arrow::float64()(对应C++double),类型不兼容会导致数组构建失败,最终Table无有效数据。 - 字段与数组数量不匹配:schema中每个行索引对应
x_i和y_i两个字段,总计2*rows个字段,但原代码仅向array_vector添加了rows个x列的数组,字段数不匹配导致Table无法正确初始化。 - 未确保输出流的正确刷新/关闭:写入操作完成后未显式刷新或关闭输出流,可能导致数据未写入磁盘。
修正后的完整代码
#include <vector> #include <arrow/api.h> #include <parquet/arrow/writer.h> int main() { // 示例初始化rows和cols,实际业务中替换为真实值 int rows = 2; int cols = 3; std::vector<std::vector<double>> x(rows, std::vector<double>(cols, 1.0)); std::vector<std::vector<double>> y(rows, std::vector<double>(cols, 2.0)); // 1. 构建匹配目标格式的Schema arrow::FieldVector fields; for(int i = 0; i < rows; i++) { fields.push_back(arrow::field("x_" + std::to_string(i), arrow::float64())); fields.push_back(arrow::field("y_" + std::to_string(i), arrow::float64())); } std::shared_ptr<arrow::Schema> schema = arrow::schema(fields); // 2. 构建包含所有x、y列的ArrayVector arrow::ArrayVector array_vector; for(int i = 0; i < rows; i++) { // 处理x列:使用DoubleBuilder匹配float64类型 arrow::DoubleBuilder x_builder; std::shared_ptr<arrow::Array> x_array; PARQUET_THROW_NOT_OK(x_builder.Reserve(cols)); for(int j = 0; j < cols; j++) { PARQUET_THROW_NOT_OK(x_builder.Append(x[i][j])); } PARQUET_THROW_NOT_OK(x_builder.Finish(&x_array)); array_vector.push_back(x_array); // 处理y列 arrow::DoubleBuilder y_builder; std::shared_ptr<arrow::Array> y_array; PARQUET_THROW_NOT_OK(y_builder.Reserve(cols)); for(int j = 0; j < cols; j++) { PARQUET_THROW_NOT_OK(y_builder.Append(y[i][j])); } PARQUET_THROW_NOT_OK(y_builder.Finish(&y_array)); array_vector.push_back(y_array); } // 3. 创建并验证Table有效性 std::shared_ptr<arrow::Table> table = arrow::Table::Make(schema, array_vector); if (!table->Validate().ok()) { return 1; } // 4. 写入Parquet文件并确保流正确关闭 std::shared_ptr<arrow::io::FileOutputStream> outfile; PARQUET_THROW_NOT_OK(arrow::io::FileOutputStream::Open("test.parquet", &outfile)); // 最后一个参数设置为表的行数,避免小数据量分片问题 PARQUET_THROW_NOT_OK(parquet::arrow::WriteTable(*table, arrow::default_memory_pool(), outfile, cols)); // 显式关闭输出流,确保缓存数据写入磁盘 PARQUET_THROW_NOT_OK(outfile->Close()); return 0; }
关键修改说明
- 类型对齐:将
FloatBuilder替换为DoubleBuilder,与schema的float64类型完全匹配。 - 补充y列数据:在循环中同时处理x和y的数组,保证
array_vector元素数量(2*rows)与schema字段数一致。 - 增强错误检查:对构建器的
Reserve、Append、Finish操作添加错误校验,提前定位问题。 - 强制流关闭:写入完成后调用
outfile->Close(),确保所有缓存数据持久化到磁盘。 - 优化写入批次:
WriteTable的批次大小设为表的总行数,避免小数据场景下的无效分片。
验证方法(Python读取)
使用pyarrow读取生成的文件,确认数据格式符合预期:
import pyarrow.parquet as pq table = pq.read_table("test.parquet") print(table.to_pandas())
输出将与需求中的表格格式一致,每行对应原二维数组的列索引,每列对应x_i或y_i的所有值。
内容的提问来源于stack exchange,提问作者Merlin486
相关产品推荐
相关产品推荐

