クエリ結果をバッチ処理する
チームは、デジタルチェックアウトの行をまとめて1つの大きな DataFrame として収集するのではなく、小さなチャンクに分けて下流のシステムに送りたいと考えています。5,000 行ずつのバッチを反復処理し、各バッチの簡単なサマリーを取得しましょう。
フィルタリング済みのデジタルチェックアウトを含む LazyFrame digital_rows はあらかじめ読み込まれています。
この演習はコースの一部です
Polars によるデータパイプラインのスケーリングと最適化
演習の手順
- ストリーミングエンジンを使って、
digital_rowsを 5,000 行ずつのバッチで反復処理してください。
実践的なインタラクティブ演習
このサンプルコードを完成させて、この演習に挑戦してみましょう。
batch_summaries = []
# Stream digital_rows in chunks of 5,000
for batch_no, batch in enumerate(
digital_rows.____(chunk_size=____, engine="streaming"),
start=1,
):
batch_summaries.append(
{
"batch": batch_no,
"rows": batch.height,
"checkouts": batch["checkouts"].sum(),
}
)
result = pl.DataFrame(batch_summaries)
print(result)