如何在Rust中使用Polars 0.53.0新流式引擎正确实现惰性扫描与Sink输出?
如何在Rust中使用Polars 0.53.0新流式引擎正确实现惰性扫描与Sink输出?
我帮你排查了代码里的问题,核心是流式模式没有在Sink阶段被正确启用,调整几个小细节后就能正常输出结果了。
问题分析
你的代码已经接近正确,但有个关键遗漏:虽然你给LazyFrame调用了with_new_streaming(true),但UnifiedSinkArgs的默认配置是关闭流式处理的,导致整个流程没有真正以流式方式执行,最终没有生成输出文件。
修正后的完整代码
Cargo.toml(保持你的配置即可)
[dependencies] polars = {version = "0.53.0", features = ["lazy", "new_streaming", "polars-io", "csv"]}
main.rs
use polars::prelude::*; use std::sync::Arc; fn process_data(intsv: &str, outtsv: &str) -> std::io::Result<()> { // 定义输入数据的Schema let schema = Schema::from_iter(vec![ Field::new("idx", DataType::Int64), Field::new("f1", DataType::Int64), Field::new("f2", DataType::Int64), ]); // 构建惰性CSV读取器,保持流式扫描 let lf1 = LazyCsvReader::new(PlRefPath::new(intsv)) .with_has_header(true) .with_separator(b'\t') .with_schema(Some(Arc::new(schema))) .finish() .expect("无法创建惰性CSV读取器"); // 配置输出TSV的格式选项 let mut seropts = SerializeOptions::default(); seropts.separator = b'\t'; let mut csvwo = CsvWriterOptions::default(); csvwo.include_header = true; csvwo.serialize_options = Arc::new(seropts); // 执行流式过滤并Sink输出,关键是开启流式Sink lf1 .with_new_streaming(true) // 启用新流式引擎 .filter(col("idx").eq(lit(1))) // 过滤idx=1的行 .sink( SinkDestination::File { target: SinkTarget::Path(PlRefPath::new(outtsv)) }, FileWriteFormat::Csv(csvwo), UnifiedSinkArgs { streaming: true, // 核心:显式开启流式输出 ..Default::default() } ) .expect("写入输出文件失败"); 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!("错误: {err:?}"); 1 } }); }
关键修改点
- 开启流式Sink:在
UnifiedSinkArgs中显式设置streaming: true,这是让整个流程真正以流式方式执行的核心; - 优化错误提示:把
unwrap()替换为expect()并添加中文提示,方便调试时定位问题; - 代码简化:移除了不必要的中间变量,让逻辑更清晰。
测试结果
运行后,mydata-filt.tsv会生成符合预期的输出:
idx f1 f2 1 4 5 1 9 3
备注:内容来源于stack exchange,提问作者Parsimonious Don
相关产品推荐
相关产品推荐

