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

Rust Polars中scan_parquet_files阻塞Tokio主线程问题求助

问题

在Azure环境中使用Polars的scan_parquet_files读取多个Parquet文件时,多数情况运行正常,但有5%-10%的概率触发panic错误:

thread 'main' panicked at 'Cannot start a runtime from within a runtime. This happens because a function (like block_on) attempted to block the current thread while the thread is being used to drive asynchronous tasks.', /home/xyx/.cargo/registry/src/github.com-1ecc6299db9ec823/tokio-0.2.21/src/runtime/enter.rs:38:5
note: run with RUST_BACKTRACE=1 environment variable to display a backtrace.

本质是Polars线程阻塞了Tokio异步主线程,尝试用spawn_blocking解决但并非100%有效,求更可靠的实现方式避免该错误。

相关代码

//Function for reading files from cloud
fn fetch_data (store: Vec<&str>,mpath:&str,args:ScanArgsParquet)-> Result<LazyFrame,PolarsError> {

let mut paths:Vec<PathBuf> = Vec::new();
for &s in store.iter(){
    let p = format!("{}{}{}{}","azure:/",mpath,&s,".parquet");
    let p1 = PathBuf::from(p);
    paths.push(p1);
}
let df = LazyFrame::scan_parquet_files(paths.into(),args).unwrap().with_streaming(true);

match df {
    a => return Ok(a),
    _  => return Err(PolarsError::NoData("Not data found".into()))
};

}

#[tokio::main]
async fn main() -> PolarsResult<()>{
//Provide Blob Storage Name and Key
let AccountName = std::env::var("STORAGE_ACC_NAME").expect("missing STORAGE_ACCOUNT");
let AccountKey = std::env::var("access_key").expect("missing STORAGE_ACCOUNT_KEY");
let ContName = std::env::var("container_name").expect("missing CONTAINER_NAME");

const TEST_S2: &str = "https://cprxt.blob.core.windows.net/blob_1";

//set up cloud options    
let cloud_options=cloud::CloudOptions::default().with_azure([(Key::AccountName,AccountName),(Key::AccessKey,AccountKey),(Key::ContainerName,ContName),(Key::Endpoint,TEST_S2)]);
let mut args = ScanArgsParquet::default();

//Set Options for scan parquet 
args.row_count=None;
args.n_rows = None;
args.low_memory=false; 
args.use_statistics = false;
args.cache= false;
args.parallel = ParallelStrategy::RowGroups;
args.cloud_options = Some(cloud_options);

//Define Master Paths 
let mpath_fp = "/read_parquet/fp_data/";
let mpath_product = "/read_parquet/product_data/";   
let mpath_sales = "/read_parquet/sales/";    
let mpath_attribute_product = "/read_parquet/attribute_product/";
let mpath_vars = "/read_parquet/vars/";

//Parquet Partitions for all the files
let mut input1 = vec!["1 ------- 33","4 ------- 33", "6 ------- 33", "7 ------- 33", "10 ------- 33"]; 

//////New function for parquet retrieval # Using scan parquet inside spawn_blocking to avoid blocking the main thread by polars
let dfs= tokio::task::spawn_blocking( move || {
    let fp = fetch_data(input1.clone(),mpath_fp,args.clone()).unwrap();
    let df_product = fetch_data(input1.clone(),mpath_product,args.clone()).unwrap();
    let df_sales = fetch_data(input1.clone(),mpath_sales,args.clone()).unwrap();
    let mut df_vars = fetch_data(input1.clone(),mpath_vars,args.clone()).unwrap();
    
    let dp= collect_all(vec![fp,df_product,df_sales,df_vars]).unwrap();
    (dp)
}).await.expect("Task panicked");

Ok(())
}
解决方案

1. 彻底隔离Polars与Tokio异步上下文

spawn_blocking失效的核心原因是Polars内部操作可能意外进入Tokio runtime上下文,需确保所有Polars相关逻辑完全脱离异步线程池:

  • 提前克隆所有非异步依赖数据,避免在spawn_blocking闭包中处理复杂克隆操作;
  • 确保闭包内仅执行Polars同步API,不调用任何Tokio异步方法。

修改后的核心代码片段:

// 提前克隆所有需要的非异步数据
let input1_clone = input1.clone();
let mpath_fp_clone = mpath_fp.to_string();
let mpath_product_clone = mpath_product.to_string();
let mpath_sales_clone = mpath_sales.to_string();
let mpath_vars_clone = mpath_vars.to_string();
let args_clone = args.clone();

let dfs = tokio::task::spawn_blocking(move || {
    // 所有Polars操作完全在阻塞线程中执行
    let fp = fetch_data(input1_clone.clone(), &mpath_fp_clone, args_clone.clone()).unwrap();
    let df_product = fetch_data(input1_clone.clone(), &mpath_product_clone, args_clone.clone()).unwrap();
    let df_sales = fetch_data(input1_clone.clone(), &mpath_sales_clone, args_clone.clone()).unwrap();
    let df_vars = fetch_data(input1_clone, &mpath_vars_clone, args_clone).unwrap();
    
    collect_all(vec![fp, df_product, df_sales, df_vars]).unwrap()
}).await.expect("Task panicked");

2. 使用Polars异步API适配Tokio runtime

Polars提供了原生异步的Parquet扫描接口scan_parquet_files_async,直接适配Tokio环境,从根源避免阻塞冲突:

  • 将fetch_data改为异步函数,替换为异步扫描API;
  • 主函数中直接异步调用,无需spawn_blocking。

修改后的代码片段:

// 异步版本的fetch_data
async fn fetch_data(store: Vec<&str>, mpath: &str, args: ScanArgsParquet) -> Result<LazyFrame, PolarsError> {
    let mut paths: Vec<PathBuf> = Vec::new();
    for &s in store.iter() {
        let p = format!("{}{}{}{}", "azure:/", mpath, &s, ".parquet");
        paths.push(PathBuf::from(p));
    }
    // 使用异步扫描API
    let df = LazyFrame::scan_parquet_files_async(paths.into(), args).await?.with_streaming(true);
    Ok(df)
}

#[tokio::main]
async fn main() -> PolarsResult<()> {
    // ... 省略配置代码 ...

    // 直接异步调用,无需spawn_blocking
    let fp = fetch_data(input1.clone(), mpath_fp, args.clone()).await?;
    let df_product = fetch_data(input1.clone(), mpath_product, args.clone()).await?;
    let df_sales = fetch_data(input1.clone(), mpath_sales, args.clone()).await?;
    let df_vars = fetch_data(input1, mpath_vars, args).await?;
    
    let dp = collect_all(vec![fp, df_product, df_sales, df_vars]).await?;

    Ok(())
}

3. 为Polars创建独立的Tokio runtime

如果必须使用同步API,可以为Polars单独创建一个独立的Tokio runtime,彻底和主runtime隔离:

use tokio::runtime::Runtime;

let dfs = tokio::task::spawn_blocking(move || {
    // 为Polars创建独立runtime
    let polars_runtime = Runtime::new().unwrap();
    polars_runtime.block_on(async {
        let fp = fetch_data(input1.clone(), mpath_fp, args.clone()).unwrap();
        let df_product = fetch_data(input1.clone(), mpath_product, args.clone()).unwrap();
        let df_sales = fetch_data(input1.clone(), mpath_sales, args.clone()).unwrap();
        let df_vars = fetch_data(input1, mpath_vars, args).unwrap();
        
        collect_all(vec![fp, df_product, df_sales, df_vars]).unwrap()
    })
}).await.expect("Task panicked");

关键注意事项

  • 避免同步/异步Polars API混用,保持上下文一致;
  • 检查cloud_options是否隐含持有Tokio runtime句柄,若有则在独立线程中初始化;
  • 升级Polars到最新稳定版,修复旧版本可能存在的异步上下文泄漏bug。

内容的提问来源于stack exchange,提问作者amit r

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 07:32:33