如何在Rust中使用新的Polars流式引擎正确实现惰性扫描和输出?
我在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_* 操作的默认选项。如果你没有输出,通常源自三种情况之一:查询在没有触发的情况下被延迟执行,静默回退到内存引擎失败,或汇聚配置不正确。
- 确保触发执行
在最新版本中,某些 sink 方法带有一个 lazy 参数,或返回一个需要最终“收集”才能真正启动数据流的句柄。如果你使用多汇聚模式,或某些方法签名已经更改为非阻塞:
- 标准汇聚(阻塞):Rust中大多数
sink_parquet调用都是即时的并返回PolarsResult<()>。确保你在处理结果(例如.expect("failed to sink"))时没有忽略。 -
惰性汇聚(非阻塞):如果你使用
collect_all或专门的流句柄,工作将在你显式轮询或收集结果之前不会开始。 -
检查静默回退
流引擎尚不支持所有操作(例如某些复杂连接或自定义UDF)。遇到不受支持的操作时,Polars可能会静默回退到 内存引擎。如果你的数据集确实“超过RAM”的大小,这种回退可能导致静默崩溃或出现看起来像“无输出”的挂起。
调试提示:用环境变量 POLARS_VERBOSE=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导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。