You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何控制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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.13 09:37:03