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 withRUST_BACKTRACE=1environment 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

