如何控制LazyFrame.map_batches流式处理模式下的小批量大小?
控制Polars LazyFrame流式处理中的mini-batch大小
在使用LazyFrame.map_batches(streamable=True)并配合collect(engine="streaming")时,mini-batch的大小可以通过以下两种方式控制:
1. 读取数据源时指定batch_size
使用scan_csv、scan_parquet等懒加载函数读取数据时,直接传入batch_size参数,该参数会决定流式处理阶段每个mini-batch的行数。示例代码:
import polars as pl # 读取CSV时指定每个batch为1000行 lf = pl.scan_csv("data.csv", batch_size=1000) # 定义可流式处理的batch逻辑 def process_batch(batch: pl.DataFrame) -> pl.DataFrame: return batch.with_columns(pl.col("num_col") * 2) # 执行流式处理并收集结果 final_result = lf.map_batches(process_batch, streamable=True).collect(engine="streaming")
2. 通过全局配置临时设置
如果不想在读取阶段单独指定,可使用polars.Config临时设置全局流式batch大小,该配置会覆盖未单独设置batch_size的流式操作:
import polars as pl # 临时将全局流式batch大小设为2000行 with pl.Config(streaming_batch_size=2000): lf = pl.scan_parquet("data.parquet") final_result = lf.map_batches(process_batch, streamable=True).collect(engine="streaming")
注意事项
map_batches本身没有直接的batch_size参数,流式模式下的batch大小由上游读取配置或全局流式配置决定。- 若同时设置了读取时的
batch_size和全局streaming_batch_size,读取阶段的参数优先级更高。
内容的提问来源于stack exchange,提问作者siminsimisim
相关产品推荐
相关产品推荐

