Tokio运行时工作线程栈溢出排查:Rust异步DataFrame操作异常
问题:WebSocket流追加Polars DataFrame时栈溢出崩溃
我编写的代码意图将WebSocket流的每个响应通过.for_each映射为AggrTrade结构体,并异步追加到Polars DataFrame df中,但程序会触发栈溢出崩溃:
- Release模式下运行到约400行时崩溃
- Debug模式下仅运行到约60行就崩溃
注释掉以下两行代码后,程序可正常运行:
append_df(&*df, df2).await; write_parquet(&*df).await;
崩溃时的错误信息:
thread 'tokio-runtime-worker' has overflowed its stack
fatal runtime error: stack overflow
[1] 29186 abort cargo run -r
2023-01-13 22:20:14.377 osascript[29344:6199784] NSNotificationCenter connection invalid
完整代码如下:
async fn new_stream(endpoint: &str) -> WebSocketStream<MaybeTlsStream<TcpStream>> { let url = "wss://fstream.binance.com/ws/".to_owned() + endpoint; let (ws_stream, _) = connect_async(&url).await.expect("Failed to connect"); ws_stream } pub async fn new_handler(endpoint: &str) -> tokio::task::JoinHandle<()> { let df = Arc::new(Mutex::new(DataFrame::empty().lazy())); let stream = new_stream(endpoint).await; let handle = tokio::spawn(async move { let df = Arc::clone(&df); stream .for_each(|msg| async { match msg { Ok(msg) => { println!("{}", &msg); let jsonmsg: Result<AggTrade, serde_json::Error> = serde_json::from_str(&msg.to_string()); let df2 = DataFrame::new(vec![jsonmsg]) .expect("Failed to create new dataframe") .lazy(); append_df(&*df, df2).await; write_parquet(&*df).await; } Err(e) => { println!("stream error: {}", e); } } }) .await }); handle } async fn append_df(df: &Mutex<LazyFrame>, df2: LazyFrame) { let mut df_edit = df.lock().await; *df_edit = concat([df_edit.clone(), df2], false, true).unwrap(); println!("{:?}", df_edit.clone().collect().unwrap()); } async fn write_parquet(df: &Mutex<LazyFrame>) { let mut df_write = df.lock().await; if &df_write.clone().collect().unwrap().height() <= &1000 { return; } println!("Writing to disk now!"); let time = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .unwrap() .as_secs(); let mut file = std::fs::File::create(&(time.to_string().trim() + ".parquet")).unwrap(); ParquetWriter::new(&mut file) .with_compression(ParquetCompression::Snappy) .with_statistics(true) .finish(&mut df_write.clone().collect().unwrap()) .unwrap(); *df_write = DataFrame::empty().lazy(); }
问题根源与修复方案
核心问题:LazyFrame的误用导致查询计划无限嵌套
崩溃的原因不是Arc<Mutex>的使用错误,而是错误地用LazyFrame存储动态追加的数据:
LazyFrame是Polars用于构建延迟执行查询计划的结构,每次调用concat都会在原有查询计划上新增一层拼接节点,而非直接合并数据- 随着数据不断追加,查询计划的嵌套深度线性增长,最终在调用
collect()时触发栈溢出(Debug模式栈空间更小,崩溃更早) - 频繁的
clone()操作复制整个查询计划,进一步加剧栈内存消耗
修复步骤
1. 切换为DataFrame存储内存数据
DataFrame适合动态追加场景,直接存储实际数据而非查询计划:
// 初始化时改为DataFrame,而非LazyFrame let df = Arc::new(Mutex::new(DataFrame::empty()));
2. 重写append_df函数
直接使用vstack合并数据,避免查询计划嵌套:
async fn append_df(df: &Mutex<DataFrame>, df2: DataFrame) { let mut df_edit = df.lock().await; // 用vstack直接合并DataFrame *df_edit = df_edit.vstack(&df2).unwrap(); println!("{:?}", df_edit); }
3. 调整new_handler中的数据转换逻辑
替换LazyFrame相关代码为DataFrame,并提前处理反序列化错误:
stream .for_each(|msg| async { match msg { Ok(msg) => { println!("{}", &msg); // 提前处理反序列化错误 let trade = match serde_json::from_str::<AggTrade>(&msg.to_string()) { Ok(t) => t, Err(e) => { println!("json parse error: {}", e); return; } }; // 直接创建DataFrame,无需转为LazyFrame let df2 = DataFrame::new(vec![trade]) .expect("Failed to create new dataframe"); append_df(&*df, df2).await; write_parquet(&*df).await; } Err(e) => { println!("stream error: {}", e); } } }) .await
4. 优化write_parquet函数
- 移除不必要的
clone()操作 - 替换为Tokio异步文件IO,避免阻塞runtime:
async fn write_parquet(df: &Mutex<DataFrame>) { let mut df_write = df.lock().await; if df_write.height() <= 1000 { return; } println!("Writing to disk now!"); let time = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .unwrap() .as_secs(); // 使用Tokio异步文件创建,避免阻塞worker线程 let mut file = tokio::fs::File::create(&(time.to_string() + ".parquet")).await.unwrap(); // 直接操作当前DataFrame,无需clone ParquetWriter::new(&mut file) .with_compression(ParquetCompression::Snappy) .with_statistics(true) .finish(&mut *df_write) .unwrap(); // 清空DataFrame *df_write = DataFrame::empty(); }
5. 可选优化
- 移除调试用的
println,减少性能开销 - 考虑批量追加数据,而非单条追加,进一步提升性能
内容的提问来源于stack exchange,提问作者ripbozo
相关产品推荐
相关产品推荐

