😊

Cloud FunctionでのETL処理:ハードコードされたファイル名を動的処理に変更する実装パターン

に公開

はじめに

GCSへのファイルアップロードをトリガーとするETLパイプラインで、処理対象ファイルがハードコーディングされていたために柔軟性を欠いていた問題を解決した記録です。

「イベント駆動型なのに特定ファイル名しか処理できない」という矛盾を、動的処理への変更で解消しました。

問題の背景

発生していた課題

Cloud Functionのmain.py内で、処理対象となるキャンペーンファイルのファイル名が以下のようにハードコーディングされていました。

# 問題のあったコード
self.campaign_files_to_load = ['運用リスト202505141.xlsx']

何が問題だったのか

このETLパイプラインは、以下のような設計でした:

GCSへファイルアップロード

Workflows トリガー

Cloud Function 実行

BigQueryへロード

イベント駆動型の設計なのに、実際には:

  • 運用リスト202505141.xlsx という特定ファイル名しか処理できない
  • 別の日付のファイル(例: 運用リスト202505142.xlsx)をアップロードしても無視される
  • ファイル名が変わるたびにコードを修正・デプロイする必要がある

これは「GCSアップロードをトリガーにする」という設計思想と矛盾していました。

解決策:3つの修正ポイント

1. 設定クラスの汎用化

まず、ETLConfigクラスで、ハードコードされた値を動的に扱えるように修正しました。

ファイルリストの初期化

# 修正前:特定ファイルをハードコード
self.campaign_files_to_load = ['運用リスト202505141.xlsx']

# 修正後:空リストで初期化(実行時に動的に追加)
self.campaign_files_to_load = []

テーブルマッピングの汎用化

BigQueryのテーブル名マッピングを、ファイル名依存から汎用キーに変更しました。

# 修正前:具体的なファイル名でマッピング
self.table_name_mapping = {
    ('運用リスト202505141.xlsx', None): 'l0_campaigns'
}

# 修正後:汎用的なキーでマッピング
self.table_name_mapping = {
    ('campaign', None): 'l0_campaigns'
}

この変更により、「キャンペーン系のファイルなら全てl0_campaignsテーブルに格納」という汎用的なルールになりました。

2. ファイル準備処理の動的化

GCSからファイルをダウンロードする_prepare_local_files_from_gcs関数で、パスに基づいて動的に処理対象を判定するロジックを追加しました。

def _prepare_local_files_from_gcs(bucket: str, name: str):
    """
    GCSから指定されたファイルをローカルにダウンロードし、
    パスに応じて動的に処理対象リストに追加する
    
    Args:
        bucket: GCSバケット名
        name: GCS内のファイルパス(例: "client_data/campaigns/report_20250524.xlsx")
    """
    storage_client = storage.Client()
    bucket_obj = storage_client.bucket(bucket)
    blob = bucket_obj.blob(name)
    
    # ローカルの一時ディレクトリにダウンロード
    local_path = f"/tmp/{os.path.basename(name)}"
    blob.download_to_filename(local_path)
    logger.info(f"Downloaded {name} to {local_path}")
    
    # 修正前:ハードコードされたファイル名と一致するかチェック
    # for i, v in enumerate(config.campaign_files_to_load):
    #     if v == file_basename:
    #         config.campaign_files_to_load[i] = local_path
    
    # 修正後:パスに基づいて動的に判定
    if "/campaigns/" in name:
        config.campaign_files_to_load.append(local_path)
        logger.info(f"Added {local_path} to campaign processing queue")

ポイント

  • ファイル名ではなくGCSパスで判定
  • /campaigns/フォルダ配下のファイルなら自動で処理対象に追加
  • ファイル名が何であろうと、配置場所で処理内容が決まる

3. BigQueryアップロード処理の修正

BigQueryへのアップロード時のテーブル特定ロジックも、汎用キーを使うように変更しました。

def upload_campaign_dataframe(df, file_path):
    """
    キャンペーンデータをBigQueryにアップロードする
    
    Args:
        df: アップロード対象のDataFrame
        file_path: 処理したファイルのローカルパス
    """
    # 修正前:ファイル名でテーブルIDを検索
    # file_basename = os.path.basename(file_path)
    # table_id = config.table_name_mapping.get((file_basename, None))
    
    # 修正後:汎用キーでテーブルIDを検索
    table_id = config.table_name_mapping.get(('campaign', None))
    
    if not table_id:
        logger.error("Campaign table mapping not found")
        raise ValueError("Table mapping for campaign not configured")
    
    # BigQueryへアップロード
    client = bigquery.Client()
    job = client.load_table_from_dataframe(df, table_id)
    job.result()
    logger.info(f"Uploaded {len(df)} rows to {table_id}")

修正後の動作フロー

GCSアップロード: gs://bucket/client_data/campaigns/任意のファイル名.xlsx

Workflows トリガー

Cloud Function起動

_prepare_local_files_from_gcs()
    ├─ パスに "/campaigns/" が含まれるか判定
    ├─ 含まれていれば処理対象リストに追加
    └─ ローカルにダウンロード

ETL処理実行
    ├─ ファイルを読み込み・変換
    └─ 汎用キー 'campaign' でテーブルを特定

BigQuery l0_campaigns テーブルへロード

この実装パターンのメリット

1. 真のイベント駆動型

ファイル名に関係なく、配置場所だけで処理内容が決まるため、運用が直感的になります。

# どのファイル名でもOK
gs://bucket/client_data/campaigns/report_20250524.xlsx
gs://bucket/client_data/campaigns/週次レポート.xlsx
gs://bucket/client_data/campaigns/臨時_緊急.xlsx

2. コード変更不要

新しいファイルをアップロードするたびにコードを修正・デプロイする必要がありません。

3. マルチカテゴリ対応

同じパターンで複数のカテゴリに対応できます。

# パスベースのルーティング
if "/campaigns/" in name:
    config.campaign_files_to_load.append(local_path)
elif "/products/" in name:
    config.product_files_to_load.append(local_path)
elif "/stores/" in name:
    config.store_files_to_load.append(local_path)

4. テーブルマッピングの管理が容易

ファイル名とテーブル名の対応を、シンプルなキーで管理できます。

table_name_mapping = {
    ('campaign', None): 'l0_campaigns',
    ('product', None): 'l0_products',
    ('store', None): 'l0_stores'
}

実装時の注意点

パス構造の設計

GCS上のディレクトリ構造を、処理カテゴリと対応させる設計にすることが重要です。

推奨構造:
gs://bucket/
  ├── client_data/
  │   ├── campaigns/    # キャンペーン系ファイル
  │   ├── products/     # 商品系ファイル
  │   └── stores/       # 店舗系ファイル

エラーハンドリング

想定外のパスにファイルが配置された場合の処理を明確にしましょう。

if "/campaigns/" in name:
    config.campaign_files_to_load.append(local_path)
elif "/products/" in name:
    config.product_files_to_load.append(local_path)
else:
    # 想定外のパス
    logger.warning(f"Unknown category in path: {name}")
    # オプション1: エラーとして扱う
    # raise ValueError(f"Unknown file category: {name}")
    # オプション2: スキップして続行
    return

テーブルマッピングの存在確認

必ずテーブルIDの取得時に存在チェックを入れましょう。

table_id = config.table_name_mapping.get(('campaign', None))
if not table_id:
    raise ValueError("Table mapping not found for 'campaign'")

まとめ

ハードコードされたファイル名を動的処理に変更することで:

  • 柔軟性: どんなファイル名でも処理可能
  • 保守性: コード変更・デプロイが不要
  • 拡張性: 新カテゴリの追加が容易
  • 整合性: イベント駆動型の設計思想と一致

特にGCSトリガーを使ったETLパイプラインでは、「配置場所で処理内容を決める」パターンが有効です。

ファイル名のルールを運用側に強制するのではなく、システム側が柔軟に対応できる設計を心がけましょう!

Discussion