如何降低Parquet文件并行批量读取时的内存占用
问题根源与解决方案
你的内存持续增长核心原因是调用RecordBatchReader->ToTable()会一次性将整个行组批次的所有数据加载到内存并生成Table对象,而Table及其内部的ChunkedArray、Array会持有大量内存引用,直到这些对象被完全销毁才会释放内存。Close()方法仅负责关闭底层IO流,不会回收已加载到内存的数据,所以调用它无法解决内存增长问题。
核心修改方案
1. 替换ToTable()为逐批次读取RecordBatch
不要一次性加载整个批次到Table,而是循环调用ReadNext()逐个读取小批次数据,处理完每个批次后立即让其脱离作用域,触发内存回收。
2. 用局部作用域包裹数据处理逻辑
将每个RecordBatch的处理放在单独代码块中,确保处理完成后,RecordBatch和衍生的Array对象立即被销毁,内存池可以及时回收内存。
3. 避免不必要的对象持有
不要将所有列的Array长期保存到容器中(除非业务必须),处理完每列数据后直接释放引用。
修改后的关键代码片段
auto read_recordbatch = [ncolumns, thread_start](size_t i, std::shared_ptr<::arrow::RecordBatchReader> reader) -> ::arrow::Result<bool>{ auto io_start = std::chrono::high_resolution_clock::now(); std::shared_ptr<arrow::RecordBatch> batch; while (ARROW_TRY(reader->ReadNext(&batch)) && batch != nullptr) { // 局部作用域包裹批次处理,确保处理完立即释放内存 { std::vector<std::shared_ptr<::arrow::Array>> vec_array; for (int col_idx = 0; col_idx < ncolumns; col_idx++) { auto array = batch->column(col_idx); // 在这里执行你的业务逻辑,比如数据解析、计算等 vec_array.emplace_back(array); } // 作用域结束,vec_array和内部的Array引用被销毁,内存可被回收 } // 旧的batch会被新批次覆盖,引用计数归零后自动释放内存 } auto io_end = std::chrono::high_resolution_clock::now(); ARROW_RETURN_NOT_OK(reader->Close()); std::cout << "thread " << i << " processed " << "io_overhead " << std::chrono::duration<double, std::milli>(io_end-io_start).count() << " ms" << std::endl; return true; };
额外优化建议
- 调整内存池闲置内存:可以在批次处理完成后调用
arrow::default_memory_pool()->ReleaseUnused(),手动释放内存池中的闲置内存(内存池默认会自动管理,此操作作为补充)。 - 控制并行度:如果CPU线程池线程数过多,同时加载的行组批次会导致内存暴涨,可根据机器内存情况调整线程池大小,或限制同时处理的批次数量。
- 取消不必要的数组合并:原代码中
ChunkedArrayToArray会将分块数组合并为单个大数组,进一步增加内存占用,若业务不需要合并,直接处理分块数组即可。
内容的提问来源于stack exchange,提问作者fcqiao
相关产品推荐
相关产品推荐

