使用Polars DataFrame写入长行时触发运行时错误
异步Rust写入Parquet大行数时运行时错误解决
问题情况
异步Rust代码作为大型程序的一部分运行,当DataFrame行数为10或30时可成功写入Parquet文件,但行数增至300时,每个异步线程都会抛出运行时错误。
DataFrame示例
df: Ok(shape: (300, 4) ┌───────────────┬─────────┬────────────────────────────────┬───────────────────────────────────────┐ │ timestamp ┆ ticker ┆ bids ┆ asks │ │ --- ┆ --- ┆ --- ┆ --- │ │ i64 ┆ str ┆ list[list[f64]] ┆ list[list[f64]] │ ╞═══════════════╪═════════╪════════════════════════════════╪═══════════════════════════════════════╡ │ 1674962575119 ┆ ETHUSDT ┆ [[1589.51, 4.731], [1589.31, ┆ [[1590.93, 39.234], [1592.1, 51.... │ │ ┆ ┆ 93.... ┆ │ │ 1674962575220 ┆ ETHUSDT ┆ [[1589.51, 22.094], [1589.31, ┆ [[1590.93, 39.234], [1592.1, 51.... │ │ ┆ ┆ 24... ┆ │ │ 1674962575319 ┆ ETHUSDT ┆ [[1589.51, 12.324], [1589.31, ┆ [[1590.93, 39.309], [1592.1, 52.... │ │ ┆ ┆ 24... ┆ │ │ 1674962575421 ┆ ETHUSDT ┆ [[1589.51, 0.0], [1589.31, ┆ [[1590.93, 26.735], [1592.1, 52.... │ │ ┆ ┆ 24.26... ┆ │ │ ... ┆ ... ┆ ... ┆ ... │ │ 1674962604998 ┆ ETHUSDT ┆ [[1440.0, 5138.446], [1558.38, ┆ [[1617.28, 40.969], [1593.72, 3.... │ │ ┆ ┆ 0... ┆ │ │ 1674962605101 ┆ ETHUSDT ┆ [[1440.0, 5138.446], [1558.38, ┆ [[1617.28, 40.969], [1593.72, 3.... │ │ ┆ ┆ 0... ┆ │ │ 1674962605201 ┆ ETHUSDT ┆ [[1440.0, 5138.446], [1558.38, ┆ [[1617.28, 40.969], [1593.72, 3.... │ │ ┆ ┆ 0... ┆ │ │ 1674962605301 ┆ ETHUSDT ┆ [[1440.0, 5138.446], [1558.38, ┆ [[1617.28, 40.969], [1593.72, 3.... │ │ ┆ ┆ 0... ┆ │ └───────────────┴─────────┴────────────────────────────────┴───────────────────────────────────────┘)
错误信息
thread '<unnamed>' panicked at 'range end index 131373 out of range for slice of length 301', C:\Users\username\.cargo\git\checkouts\arrow2-945af624853845da\baa2618\src\io\parquet\write\mod.rs:171:37 thread '<unnamed>' panicked at 'range end index 131373 out of range for slice of length 301', C:\Users\username\.cargo\git\checkouts\arrow2-945af624853845da\baa2618\src\io\parquet\write\mod.rs:171:37
相关代码
async { if main_vec.length() >= ROWS { let df = main_vec.to_df(); let time = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .unwrap() .as_secs(); let file = std::fs::File::create(&(time.to_string().trim() + ".parquet")).unwrap(); println!("Wrote parquet file: {}", &(time.to_string().trim().to_owned() + ".parquet")); // ERROR OCCURS HERE ParquetWriter::new(file) .with_compression(ParquetCompression::Snappy) .with_statistics(true) .finish(&mut df.collect().unwrap()) .expect("Failed to write parquet file"); main_vec.clear(); } }.await;
解决方案
1. 修复数据竞争问题
异步环境下main_vec可能被多线程同时访问,导致生成的DataFrame结构损坏。用线程安全容器包裹main_vec:
use std::sync::{Arc, Mutex}; // 初始化时 let main_vec = Arc::new(Mutex::new(YourVecType::new())); // 异步块中访问 let mut vec_guard = main_vec.lock().unwrap(); if vec_guard.length() >= ROWS { let df = vec_guard.to_df(); // ... 后续写入逻辑 vec_guard.clear(); }
2. 升级arrow2依赖
报错来自arrow2旧版本的已知bug,更新到最新稳定版即可修复:
# Cargo.toml中更新依赖 arrow2 = "0.18" # 替换为当前最新稳定版本号
3. 临时关闭统计信息写入
嵌套类型的统计信息计算可能触发越界,先尝试移除.with_statistics(true):
ParquetWriter::new(file) .with_compression(ParquetCompression::Snappy) .finish(&mut df.collect().unwrap()) .expect("Failed to write parquet file");
4. 验证DataFrame转换结果
写入前确认转换后的记录批行数正确,排除数据损坏:
let mut record_batch = df.collect().unwrap(); println!("Record batch rows: {}", record_batch.num_rows()); // 应输出300
5. 改用异步文件操作(可选)
异步环境下使用同步文件IO可能引发潜在问题,基于tokio runtime时改用异步文件:
use tokio::fs::File; // ... let file = File::create(&(time.to_string().trim() + ".parquet")).await.unwrap(); // 若ParquetWriter不支持异步,可使用tokio::io::BufWriter包装
内容的提问来源于stack exchange,提问作者ripbozo
相关产品推荐
相关产品推荐

