如何在Rust中使用新的Polars流式引擎正确实现惰性扫描和输出?

编程语言 2026-07-12

我在Rust中使用Polars 0.53.0。据我所知,使用新的流式引擎对惰性扫描和汇聚(不进行collect() / 将数据物化)的方法如下,但我没有输出。我找不到一个实际用0.53.0实现此功能的代码示例。有人知道怎么做吗?

Cargo.toml

[dependencies]
polars = {version = "0.53.0", features = ["lazy", "new_streaming", "polars-io", "csv"]}

main.rs

use polars::prelude::*;

fn process_data(intsv: &str, outtsv: &str) -> std::io::Result<()> {
    let schema = Schema::from_iter(vec![
        Field::new("idx".into(), DataType::Int64),
        Field::new("f1".into(), DataType::Int64),
        Field::new("f2".into(), DataType::Int64),
    ]);
    let lf1: LazyFrame =
        LazyCsvReader::new(PlRefPath::new(&*intsv))
        .with_has_header(true)
        .with_separator(b'\t')
        .with_schema(Some(Arc::new(schema.clone())))
        .finish().unwrap();
    let outpath = PlRefPath::new(&*outtsv);
    let mut seropts = SerializeOptions::default();
    seropts.separator = b'\t';
    let aseropts = Arc::new(seropts);
    let mut csvwo = CsvWriterOptions::default();
    csvwo.include_header = true;
    csvwo.serialize_options = aseropts;
    let frame = lf1
        .with_new_streaming(true)
        .filter(col("idx").eq(lit(1)))
        //.collect()
        .sink(
            SinkDestination::File {
                target: SinkTarget::Path(outpath)
            },
            FileWriteFormat::Csv(csvwo),
            UnifiedSinkArgs::default())
        .expect("problem writing CSV output");
    //println!("{}", frame);
    Ok(())
}

fn main() {
    let input_tsv = "mydata.tsv";
    let output_tsv = "mydata-filt.tsv";
    std::process::exit(match process_data(input_tsv, output_tsv) {
        Ok(_) => 0,
        Err(err) => {
            eprintln!("error: {err:?}");
            1
        }
    });
}

mydata.tsv

idx     f1      f2
1       4       5
2       6       7
3       7       8
1       9       3
6       0       3
5       2       3

解决方案

在Polars 0.53.0 中,“新的”流式引擎被设计为对 sink_* 操作的默认选项。如果你没有输出,通常源自三种情况之一:查询在没有触发的情况下被延迟执行,静默回退到内存引擎失败,或汇聚配置不正确。

  1. 确保触发执行

在最新版本中,某些 sink 方法带有一个 lazy 参数,或返回一个需要最终“收集”才能真正启动数据流的句柄。如果你使用多汇聚模式,或某些方法签名已经更改为非阻塞:

  • 标准汇聚(阻塞):Rust中大多数 sink_parquet 调用都是即时的并返回 PolarsResult<()>。确保你在处理结果(例如 .expect("failed to sink"))时没有忽略。
  • 惰性汇聚(非阻塞):如果你使用 collect_all 或专门的流句柄,工作将在你显式轮询或收集结果之前不会开始。

  • 检查静默回退

流引擎尚不支持所有操作(例如某些复杂连接或自定义UDF)。遇到不受支持的操作时,Polars可能会静默回退到 内存引擎。如果你的数据集确实“超过RAM”的大小,这种回退可能导致静默崩溃或出现看起来像“无输出”的挂起。

调试提示:用环境变量 POLARS_VERBOSE=1 运行你的二进制程序。这将打印控制台,指明是否实际使用了流引擎,还是回退到默认引擎。

  1. 实现示例(Polars 0.53.0)

为了在Rust中正确使用流引擎且不在RAM中把数据物化,请使用 scan_parquet(或 scan_csv),然后再跟随 sink_parquet

rust

use polars::prelude::*;

fn main() -> PolarsResult<()> {
    // 1. Point to your source (must be a 'scan' to enable streaming)
    let lf = LazyFrame::scan_parquet("input_big_data.parquet", Default::default())?
        .filter(col("some_column").gt(lit(50))) // Example transformation
        .with_streaming(true); // Explicitly hint to use streaming engine

    // 2. Sink directly to disk
    // This triggers the streaming execution immediately
    lf.sink_parquet(
        "output_result.parquet",
        ParquetWriteOptions::default(),
    )?;

    println!("Streaming complete!");
    Ok(())
站内所有文章版权归属LeftHeroAI导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。

相关文章