0Pricing
Pandas & NumPy Academy · レッスン

チャンク間の逐次集計

ファイル全体をメモリに保持せず、チャンクごとに累積件数、合計、最小値・最大値を更新します。

「チャンク間の逐次集計」はCoddyKit上の無料Pandas & NumPy Academyレッスンです。 これはレッスン2/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはPandas & NumPy Academy学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 Pandas & NumPy Academyコースには全4レッスンが含まれています。

なぜ段階的な集計を行うのか

段階的な集計は、複数のマシンに計算を分散させずに、RAMより大きなデータセットを分析するための鍵です。最終的な統計量を計算するためにすべてのデータを読み込むのではなく、各チャンクで更新する累積値(部分合計、件数、最小値・最大値)を保持します。ファイル全体を走査した後、これらの軽量な累積値から最終結果を組み立てます。このパターンを使えば、1台のノートパソコンでもテラバイト規模のファイルを処理できます。

行数のカウントと平均値の計算

チャンク全体の平均値を計算するには、累積合計と件数を別々に追跡する必要があります。チャンクごとの平均値を単純に平均してはいけません。チャンクごとにサイズが異なる可能性があるためです。正しい式は total_sum / total_count です。このパターンは、分散、相関、ヒストグラムなど、分解可能なあらゆる量に応用できます。

import pandas as pd

total_sum = 0.0
total_count = 0

for chunk in pd.read_csv('transactions.csv', chunksize=100000):
    total_sum += chunk['amount'].sum()
    total_count += chunk['amount'].notna().sum()

grand_mean = total_sum / total_count
print(f'Rows processed: {total_count:,}')
print(f'Grand mean: {grand_mean:.4f}')

最小値と最大値の段階的な計算

チャンク全体の最小値と最大値を追跡するのは簡単です。Pythonの float('inf') と float('-inf') で初期化し、各チャンクの最小値・最大値で更新します。これにより中間データを保存する必要がなくなります。グループ識別子をキーとする辞書を保持すれば、グループごとの最小値・最大値にも応用できます。

import pandas as pd

global_min = float('inf')
global_max = float('-inf')

for chunk in pd.read_csv('prices.csv', chunksize=50000):
    chunk_min = chunk['price'].min()
    chunk_max = chunk['price'].max()
    if chunk_min < global_min:
        global_min = chunk_min
    if chunk_max > global_max:
        global_max = chunk_max

print(f'Price range: {global_min} to {global_max}')

度数の段階的なカウント

カテゴリ列では、各チャンクの value_counts() の結果をPandas Seriesの累積値に加算して、累積度数辞書を維持します。Pandas Seriesの加算ではインデックスラベルが揃えられるため、後のチャンクに現れる未知のカテゴリも自動的に含まれます。すべてのチャンクを処理した後、件数順に並べ替えると、データセット全体で上位のカテゴリを確認できます。

import pandas as pd

freq = pd.Series(dtype='int64')

for chunk in pd.read_csv('orders.csv',
                         chunksize=100000,
                         usecols=['category']):
    chunk_counts = chunk['category'].value_counts()
    freq = freq.add(chunk_counts, fill_value=0)

# Final sorted frequency table
print(freq.sort_values(ascending=False).head(10))

GroupBy集計の段階的な実行

チャンク全体でgroupbyの合計または件数を計算するには、各チャンク内で groupby().agg() を適用し、結果のSeriesまたはDataFrameを保存します。ループの後にすべての部分結果を連結し、2回目のgroupbyを適用して結合します。この2段階の方法では、複数のチャンクに現れるグループも正しく処理できます。データがグループ別ではなく日付順に並んでいる場合に、こうした状況がよく発生します。

import pandas as pd

partials = []
for chunk in pd.read_csv('sales.csv',
                         chunksize=100000,
                         usecols=['region', 'product', 'revenue']):
    p = chunk.groupby(['region', 'product'])['revenue'].sum()
    partials.append(p)

final = (
    pd.concat(partials)
    .groupby(level=['region', 'product'])
    .sum()
    .sort_values(ascending=False)
)
print(final.head(10))

分散の段階的な計算(Welford法)

チャンク全体の分散を計算するのは、平均値より難しくなります。単純な式 E[X²] - E[X]² は、平均値が大きい場合に桁落ちを起こします。Welfordのオンラインアルゴリズムは、累積平均と偏差平方和を保持し、新しい値ごとに数値的に安定した方法で更新します。SciPyにも実装されていますが、このパターンを理解しておけば、重み付き分散や共分散にも応用できます。

import pandas as pd
import numpy as np

# Simple two-pass approach using stored chunk stats
chunk_stats = []
for chunk in pd.read_csv('data.csv',
                         chunksize=100000,
                         usecols=['value']):
    n = chunk['value'].count()
    mean = chunk['value'].mean()
    var = chunk['value'].var(ddof=1)
    chunk_stats.append((n, mean, var))

# Combine: use pooled variance formula
total_n = sum(s[0] for s in chunk_stats)
total_mean = sum(s[0]*s[1] for s in chunk_stats) / total_n
pooled_var = sum((s[0]-1)*s[2] + s[0]*(s[1]-total_mean)**2
                 for s in chunk_stats) / (total_n - 1)
print(f'Grand variance: {pooled_var:.4f}')

段階的なヒストグラムの作成

RAMに収まらない大きさのファイル全体で列の分布を計算するには、段階的なヒストグラムが必要です。まず、小さなサンプルまたはドメイン知識に基づいてビンの境界を固定します。次に各チャンク内で np.histogram(chunk_values, bins=edges) を使い、件数を累積します。最後に、結合した件数を棒グラフとして描画します。Kafka StreamsやFlinkなどのストリーミングシステムでは、この方法で近似ヒストグラムを計算しています。

import pandas as pd
import numpy as np

# Decide bin edges from a sample
sample = pd.read_csv('amounts.csv', nrows=5000)
bins = np.linspace(sample['amount'].min(),
                   sample['amount'].max(), 21)
counts = np.zeros(len(bins) - 1, dtype='int64')

for chunk in pd.read_csv('amounts.csv',
                         chunksize=100000,
                         usecols=['amount']):
    chunk_counts, _ = np.histogram(
        chunk['amount'].dropna(), bins=bins
    )
    counts += chunk_counts

print('Histogram counts:', counts[:5], '...')

一意な値の近似的な追跡

チャンク全体で正確な異なる値の数を数えるには、すべての一意な値を保存する必要があり、数百万件に達する可能性があります。大規模データで近似的にカウントするには、hyperloglog ライブラリを通じてPythonで利用できるHyperLogLogスケッチを使用します。別の方法として、チャンクごとにsetで一意な値を追跡して和集合を求めることもできますが、データ量に応じて際限なく増加します。低コストで近似するなら、チャンクごとに pd.Series.nunique() を使って平均値を報告する方法があります。正確ではありませんが、データプロファイリングには十分なことがよくあります。

import pandas as pd

unique_ids = set()
for chunk in pd.read_csv('events.csv',
                         chunksize=100000,
                         usecols=['user_id']):
    unique_ids.update(chunk['user_id'].dropna().unique())

print(f'Distinct user IDs: {len(unique_ids):,}')
# Warning: the set may grow large for high-cardinality columns

長時間実行中の進捗表示

数GBのファイルを処理すると、数分かかることがあります。パイプラインが実行中であることを確認し、残り時間を推定できるように、進捗表示を追加してください。処理したバイト数または行数を数え、ファイルサイズと比較します。tqdm ライブラリを使えば、tqdm(reader) ラッパーによって簡単に実現できます。tqdmを使わない場合でも、10チャンクごとにステータス行を出力すれば、長時間のバッチ処理中に役立つフィードバックを得られます。

import pandas as pd
import time

chunksize = 100000
start = time.time()
rows_processed = 0

for i, chunk in enumerate(pd.read_csv('big.csv',
                                       chunksize=chunksize)):
    rows_processed += len(chunk)
    # Report every 10 chunks
    if (i + 1) % 10 == 0:
        elapsed = time.time() - start
        rate = rows_processed / elapsed
        print(f'Chunk {i+1}: {rows_processed:,} rows '
              f'@ {rate/1000:.0f}k rows/sec')

print(f'Total: {rows_processed:,} rows in {time.time()-start:.1f}s')

集計前のフィルタリング

不要なデータの蓄積を避けるため、集計の前に各チャンク内でフィルターを適用します。たとえば2024年の注文だけが必要な場合は、groupbyの前にチャンクの日付列をフィルタリングします。これにより部分結果に必要なメモリを減らし、最終的な連結処理も高速化できます。効率的なデータ処理の基本原則として、パイプラインのできるだけ早い段階でフィルターを適用してください。

import pandas as pd

partials = []
for chunk in pd.read_csv('orders.csv',
                         chunksize=100000,
                         parse_dates=['order_date']):
    # Filter early: only 2024 orders
    mask = chunk['order_date'].dt.year == 2024
    filtered = chunk.loc[mask, ['category', 'revenue']]

    if len(filtered) > 0:
        p = filtered.groupby('category')['revenue'].sum()
        partials.append(p)

if partials:
    result = pd.concat(partials).groupby(level=0).sum()
    print(result)

中間結果の保存

非常に長時間実行されるジョブでは、プロセスが中断された場合にチェックポイントから再開できるよう、途中結果を定期的に保存してください。N個のチャンクを処理するごとに、チャンク単位の集計結果をParquetまたはCSVファイルに書き込みます。1000個中800個目のチャンクでジョブが失敗しても、保存した集計結果を読み込んで、ファイル全体を再処理するのではなく中断した箇所から続行できます。このレジリエンスパターンは、本番環境のデータパイプラインに不可欠です。

import pandas as pd
import os

CHECKPOINT = 'checkpoint.csv'
running_total = 0.0
running_count = 0

# Resume from checkpoint if it exists
if os.path.exists(CHECKPOINT):
    ckpt = pd.read_csv(CHECKPOINT)
    running_total = ckpt['total'].iloc[0]
    running_count = int(ckpt['count'].iloc[0])
    print(f'Resuming from checkpoint: {running_count:,} rows')

for chunk in pd.read_csv('huge.csv', chunksize=100000):
    running_total += chunk['value'].sum()
    running_count += len(chunk)

# Save checkpoint
pd.DataFrame({'total': [running_total],
              'count': [running_count]}).to_csv(CHECKPOINT, index=False)
print(f'Final mean: {running_total / running_count:.4f}')

理解度チェック

このレッスンで学んだデータ分析の概念について理解度を確認します。

レッスンのまとめ

このレッスンでは、実行中のアキュムレーター(sum、count、min/max、頻度Series)によって、大きなファイル全体を定数メモリで集計できること、2段階のgroupby(チャンクごとの部分的なgroupbyを行った後、結合して再グループ化する方法)によって、チャンクをまたぐグループを正しく処理できること、そして早期フィルタリングによって各チャンク内で集計処理にかかるコストを削減できることを学びました。次は、大規模データセットでPandasの代わりにそのまま使える並列処理ライブラリ、Dask DataFramesについて学びます。

よくある質問

「チャンク間の逐次集計」レッスンは無料ですか?

はい。「チャンク間の逐次集計」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、Pandas & NumPy Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 Pandas & NumPy Academyコースには全4レッスンが含まれています。

「チャンク間の逐次集計」で何を学びますか?

ファイル全体をメモリに保持せず、チャンクごとに累積件数、合計、最小値・最大値を更新します。 ブラウザで直接実行するハンズオンコードでPandas & NumPy Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。

Pandas & NumPy Academyを始めるのに経験は必要ですか?

事前経験は必要ありません。CoddyKitのPandas & NumPy Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン2/4です。

「チャンク間の逐次集計」レッスンにはどのくらい時間がかかりますか?

ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。

このPandas & NumPy Academyレッスンでコードを書いて実行できますか?

はい。すべてのPandas & NumPy Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。

このコースのすべてのレッスン

  1. chunksizeによるCSVのストリーミング
  2. チャンク間の逐次集計
  3. Dask DataFrame入門
  4. Parquet:高速な列指向ストレージ
← Pandas & NumPy Academyに戻る