インクリメンタル処理: 高速なパイプライン更新

この記事では、SEOパイプラインがインクリメンタル処理を利用して、数時間ではなく数秒で実行される仕組みを説明します。

問題: フル再処理は遅い

パイプライン全体を最初から実行するには数時間かかります:

  • ステップ 0 (ソース埋め込み): 15分 (65,000製品)

  • ステップ 1 (クエリ取得): 10分 (API呼び出し)

  • ステップ 2 (クエリクラスタリング): 30分 (65K×65K類似度計算)

  • ステップ 3 (フレーズマッピング): 20分 (埋め込み + マッチング)

  • ステップ 4 (製品マッチング): 45分 (クエリ × 製品)

  • ステップ 5 (関連検索): 25分 (クエリ × クエリ類似度)

合計: フルパイプラインで約2.5時間

問題点: 日次更新では、変更のないデータを再計算するために2.5時間を無駄にすることになります。

解決策: 3層インクリメンタル戦略

不要な作業をスキップするために3つの技術を使用しています:

1. ステップスキップ (大粒度)

出力が新しく、スクリプトが変更されていない場合、ステップ全体をスキップします。

2. インクリメンタル埋め込み (中粒度)

新規または変更されたアイテムのみを埋め込み、キャッシュされた埋め込みを再利用します。

3. チェックポイント (細粒度)

長時間の操作中に進捗を保存し、失敗時にはチェックポイントから再開します。

ステップスキップ戦略

仕組み

各ステップの前に以下を確認します:

出力は存在するか? 存在しない場合、ステップを実行します。

出力の新旧: 7日以上古い場合、ステップを実行します。

スクリプトは変更されたか? 出力が生成されてからスクリプトが変更された場合、ステップを実行します。

すべてのチェックに合格したか? ステップをスキップします。

実装

def should_skip_step(output_path, script_path, days=7):
    # 出力が存在するかチェック
    if not os.path.exists(output_path):
# ... (実装の詳細は省略)

使用方法

各スクリプトは起動時にチェックします:

from seo_common import should_skip_step

if should_skip_step(SEO_SOURCE_EMBEDDINGS_PATH, __file__):
    print("✓ スキップ: 出力は新しく、スクリプトは変更されていません")
    return

利点

高速な日次実行: データが変更されなければほとんどのステップがスキップされます

自動的な無効化: スクリプトの変更が再実行をトリガーします

設定可能な鮮度: ステップごとに days パラメータを調整できます

インクリメンタル埋め込み戦略

仕組み

アイテム(製品、クエリ、フレーズ)を埋め込む際に:

キャッシュをロード: 以前に埋め込まれたアイテムとそのキーを読み込みます

キーを比較: 新規、変更、削除されたアイテムを特定します

新規のみ埋め込み: キャッシュにないアイテムのみを埋め込みます

マージ: キャッシュされた埋め込みと新しい埋め込みを正しい順序で結合します

保存: 更新されたキャッシュを書き込みます

実装

incremental_embed_with_keys 関数がこれを処理します:

def incremental_embed_with_keys(
    items,           # 埋め込む現在のアイテム
    keys,            # アイテムの一意キー
# ... (実装の詳細は省略)

キャッシュヒット率

初回実行後の典型的なキャッシュヒット率:

ソースデータ (製品、部品、記事):

  • 初回実行: 0% (65,000アイテムすべて埋め込み)

  • 日次実行: 99.5% (約300の新規/変更アイテムのみ)

クエリ (GSC、広告、ライブから):

  • 初回実行: 0% (65,000クエリすべて埋め込み)

  • 日次実行: 99.2% (約500の新規クエリのみ)

フレーズマッピング:

  • 初回実行: 0% (5,000フレーズすべて埋め込み)

  • 日次実行: 99.8% (約10の新規フレーズのみ)

パフォーマンスへの影響

初回実行 (コールドキャッシュ):

  • ソース埋め込み: 15分 (65,000アイテム)

  • クエリ埋め込み: 10分 (65,000クエリ)

  • フレーズ埋め込み: 2分 (5,000フレーズ)

日次実行 (ウォームキャッシュ):

  • ソース埋め込み: 10秒 (300アイテム、99.5%ヒット率)

  • クエリ埋め込み: 5秒 (500クエリ、99.2%ヒット率)

  • フレーズ埋め込み: 1秒 (10フレーズ、99.8%ヒット率)

高速化: 90~180倍速

チェックポイント戦略

仕組み

長時間実行される操作(65,000アイテムの埋め込みなど)に対して:

バッチ処理: アイテムをバッチで処理します(例: 1,000アイテムごと)

チェックポイントを保存: 各バッチ後に蓄積された結果を保存します

失敗時再開: プロセスがクラッシュした場合、最後のチェックポイントから再開します

最終保存: すべてのバッチ後に完全な結果を保存します

実装

チェックポイント機能は incremental_embed_with_keys に組み込まれています:

checkpoint_every = 1000  # 1,000アイテムごとに保存

embeddings_list = []
# ... (実装の詳細は省略)

利点

クラッシュからの回復: 最初からやり直すのではなく、最後のチェックポイントから再開します

進捗の可視化: 1,000アイテムごとに進捗を確認できます

メモリ効率: バッチで処理し、すべてを一度にロードしません

パイプライン全体での統合

インクリメンタル処理は複数のステップで使用されています:

ステップ 0: ソースデータ埋め込み

インクリメンタル: 新規/変更された製品、部品、記事のみを埋め込みます

チェックポイント: 1,000アイテムごとに保存します

スキップロジック: 出力が7日未満でスクリプトが変更されていない場合スキップします

参照: ソースデータ埋め込み

ステップ 1: クエリ取得

インクリメンタル: API呼び出しは新規データのみを取得します(前回実行以降)

スキップロジック: 出力が1日未満の場合スキップします

参照: クエリ取得

ステップ 3b: クエリ埋め込み

インクリメンタル: 新規クエリのみを埋め込みます

チェックポイント: 1,000クエリごとに保存します

スキップロジック: 出力が7日未満でスクリプトが変更されていない場合スキップします

参照: クエリ埋め込み

ステップ 4: フレーズマッピング拡張

インクリメンタル: 新規フレーズのみを埋め込みます

チェックポイント: 1,000フレーズごとに保存します

スキップロジック: 出力が7日未満でスクリプトが変更されていない場合スキップします

参照: フレーズ-フィルターマッピング

ステップ 6: 製品マッチング

インクリメンタル: 新規クエリのみをマッチングします

スキップロジック: 出力が7日未満でスクリプトが変更されていない場合スキップします

参照: 製品マッチング

設定

インクリメンタル処理はステップごとに設定されます:

鮮度の閾値

# 出力が7日未満の場合スキップ (デフォルト)
should_skip_step(output_path, script_path, days=7)

# 出力が1日未満の場合スキップ (頻繁に変更されるデータ用)
should_skip_step(output_path, script_path, days=1)

チェックポイント頻度

# 1,000アイテムごとに保存 (デフォルト)
incremental_embed_with_keys(..., checkpoint_every=1000)

# 5,000アイテムごとに保存 (処理を速く、安全性を低く)
incremental_embed_with_keys(..., checkpoint_every=5000)

バッチサイズ

# バッチあたり32アイテム埋め込み (デフォルト、バランス型)
incremental_embed_with_keys(..., batch_size=32)

# バッチあたり64アイテム埋め込み (GPU上で高速、メモリ使用量多)
incremental_embed_with_keys(..., batch_size=64)

監視とデバッグ

キャッシュ統計

各ステップでキャッシュ統計が表示されます:

✓ 既存のキャッシュを発見、変更を確認中...
  既存: 65,000アイテム
  現在: 65,300アイテム
  再利用: 64,800埋め込み
  新規: 500アイテムを埋め込み

スキップメッセージ

ステップがスキップされるとき:

✓ 0_embed_source_data.py をスキップ: 出力は新しく、スクリプトは変更されていません。

チェックポイントメッセージ

長時間の操作中:

65,000アイテムを埋め込み中 (1,000ごとにチェックポイント)...
  バッチ 0-1000...
    ✓ チェックポイント保存 (合計 1,000)
  バッチ 1000-2000...
    ✓ チェックポイント保存 (合計 2,000)
  ...

参考文献

技術的概念

関連記事

まとめ

インクリメンタル処理により、パイプラインは90~180倍高速になります:

3層戦略:

  • ✅ ステップスキップ (出力が新しい場合ステップ全体をスキップ)

  • ✅ インクリメンタル埋め込み (新規/変更アイテムのみ埋め込み)

  • ✅ チェックポイント (進捗を保存、失敗時再開)

パフォーマンス:

  • ✅ 初回実行: 約2.5時間 (フルパイプライン)

  • ✅ 日次実行: 約5分 (インクリメンタル更新)

  • ✅ キャッシュヒット率: 初回実行後99%以上

利点:

  • ✅ 高速な日次更新 (時間ではなく数分で)

  • ✅ 自動的な無効化 (スクリプト変更が再実行をトリガー)

  • ✅ クラッシュからの回復 (チェックポイントから再開)

  • ✅ メモリ効率 (バッチ処理)

この戦略により、変更のないデータの計算にリソースを浪費することなく、日次パイプライン実行が可能になります。


← ドキュメントインデックスに戻る