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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 23:50:32