📚

AdMobの広告収益データを毎日自動取得するシステムを構築した話

に公開

はじめに

本年4月にWED株式会社へ新卒で入社いたしました、データエンジニアのmomokun7です。

毎日のように発生するアプリの広告収益データ集計。AdMobの管理画面から必要な数字を一つ一つコピー&ペーストして、Googleスプレッドシートに貼り付けている方、正直に手を挙げてください!
その作業、どれくらいの時間を費やしていますか?「もっと手軽にデータを見たい」「分析に時間を割きたいのに…」そんなジレンマ、ありませんか?

弊社は、この非効率な「手動データ転記」という悪習を断ち切るべく、Cloud Runを活用したAdMob広告収益データの自動収集・BigQuery格納システムを構築しました。

この記事では、これらの仕組みや具体的な方法について解説していきます。

アーキテクチャ概要

今回構築した仕組みでは、以下のようなフローで処理が行われます。

具体的な処理の流れ

  1. まず、開発者 (Developer) は、データ収集・処理を行うプログラムを含んだDocker Imageを作成し、Artifact Registryにプッシュして登録します。
  2. このプログラムの実行に際して必要となる、AdMob APIへのアクセスに必要な認証は、Secret Managerを通してセキュアに行います。
  3. 毎日、Cloud Schedulerが毎朝6時にイベントを発火し、AdMobデータ収集プログラムの実行をトリガーします。
  4. トリガーされたCloud Runは、Secret Managerから認証情報を取得し、AdMob APIを通じて収益データを取得します。
  5. 最後に、Cloud Runで取得されたAdMobデータは、BigQueryのテーブルに格納されます。

実装

認証トークンの保存

Google Cloud Secret Managerを使用してAPI認証情報やアクセストークンなどの機密情報を安全に管理する機能です。必要に応じて機密情報を取得し、新しいトークンが生成された際には自動的に保存・更新します。

secret_manager.py
import json
import logging
import pickle
import google.cloud.secretmanager as secretmanager
import src.config as config

logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')

def get_secret(secret_type: str):
    client = secretmanager.SecretManagerServiceClient()
    project_id = config.SECRET_MANAGER_PROJECT_ID

    if secret_type == "credentials":
        secret_id = config.CREDENTIALS_SECRET_MANAGER_SECRET_ID
        secret_name = f"projects/{project_id}/secrets/{secret_id}"

        try:
            client.get_secret(request={"name": secret_name})
        except secretmanager.exceptions.NotFound:
            raise ValueError(f"認証情報シークレット '{secret_id}' が見つかりません。")
        except Exception as e:
            raise ValueError(f"認証情報シークレット '{secret_id}' にアクセスできません: {e}")

        version_name = f"{secret_name}/versions/latest"
        response = client.access_secret_version(request={"name": version_name})
        payload = response.payload.data.decode("utf-8")
        return json.loads(payload)

    elif secret_type == "token":
        secret_id = config.TOKEN_SECRET_MANAGER_SECRET_ID
        secret_name = f"projects/{project_id}/secrets/{secret_id}"

        try:
            client.get_secret(request={"name": secret_name})
        except secretmanager.exceptions.NotFound:
            raise ValueError(f"トークンシークレット '{secret_id}' が見つかりません。")
        except Exception as e:
            raise ValueError(f"トークンシークレット '{secret_id}' にアクセスできません: {e}")

        version_name = f"{secret_name}/versions/latest"
        response = client.access_secret_version(request={"name": version_name})
        token = pickle.loads(response.payload.data)
        return token
    else:
        raise ValueError("無効なsecret_typeです。「credentials」または「token」である必要があります。")

def save_token(creds, update_token: bool = False):
    creds_data = pickle.dumps(creds)
    client = secretmanager.SecretManagerServiceClient()
    project_id = config.SECRET_MANAGER_PROJECT_ID
    secret_id = config.TOKEN_SECRET_MANAGER_SECRET_ID
    secret_name = f"projects/{project_id}/secrets/{secret_id}"

    try:
        client.get_secret(request={"name": secret_name})
        logging.info(f"シークレット '{secret_id}' は既に存在します。")
    except secretmanager.exceptions.NotFound:
        logging.info(f"シークレット '{secret_id}' が見つかりません。作成します...")
        client.create_secret(
            request={
                "parent": f"projects/{project_id}",
                "secret_id": secret_id,
                "secret": {"replication": {"automatic": {}}},
            }
        )
        logging.info(f"シークレット '{secret_id}' を作成しました。")
    except Exception as e:
        logging.error(f"シークレット '{secret_id}' の確認または作成に失敗しました: {e}")
        raise

    if update_token:
        try:
            versions = client.list_secret_versions(request={"parent": secret_name})
            for version in versions:
                if version.state == secretmanager.SecretVersion.State.ENABLED and \
                   version.name != f"{secret_name}/versions/latest":
                    client.destroy_secret_version(request={"name": version.name})
                    logging.info(f"古いシークレットバージョンを削除しました: {version.name}")
        except Exception as e:
            logging.warning(f"古いバージョンの削除中にエラーが発生しました: {e}")

    try:
        client.add_secret_version(
            request={
                "parent": secret_name,
                "payload": {"data": creds_data},
            }
        )
        logging.info(f"新しいトークンバージョンをSecret '{secret_id}' に保存しました。")
    except Exception as e:
        logging.error(f"新しいトークンの保存中にエラーが発生しました: {e}")

AdMobの認証を実行

Google APIの認証を自動化するスクリプトです。まず保存された既存のトークンが有効かチェックし、有効であればそのまま使用します。トークンが期限切れの場合は自動的に更新し、トークンが存在しない場合は新規でOAuth認証フローを実行してユーザー認証を行います。

auth_manager.py
import os
import logging
from googleapiclient.discovery import build
from google_auth_oauthlib.flow import InstalledAppFlow
from google.oauth2.credentials import Credentials
from google.auth.transport.requests import Request
import src.config as config
import src.secret_manager as sm

logger = logging.getLogger(__name__)

def authenticate(initial_creds_dict):
    client_id = initial_creds_dict.get('CLIENT_ID')
    client_secret = initial_creds_dict.get('CLIENT_SECRET')
    if not client_id or not client_secret:
        logger.error("CLIENT_IDまたはCLIENT_SECRETが見つかりません")
        return None
    
    creds_obj = None
    # 既存のトークンを試行
    try:
        creds_obj = sm.get_secret("token")
        if isinstance(creds_obj, dict):
            creds_obj = Credentials.from_authorized_user_info(creds_obj)
        if creds_obj and creds_obj.valid:
            logger.info("既存の有効なトークンを使用しています")
            return build(config.API_NAME, config.API_VERSION, credentials=creds_obj)
        elif creds_obj and creds_obj.expired and creds_obj.refresh_token:
            creds_obj.refresh(Request())
            sm.save_token(creds_obj, update_token=True)
            logger.info("期限切れトークンを更新しました")
            return build(config.API_NAME, config.API_VERSION, credentials=creds_obj)
    except Exception as e:
        logger.error(f"トークン認証に失敗しました: {e}")
        creds_obj = None # エラー発生時は新規認証へ
    
    # 新規認証を実行
    try:
        client_config = {
            "installed": {
                "client_id": client_id,
                "client_secret": client_secret,
                "redirect_uris": config.REDIRECT_URIS,
                "auth_uri": config.AUTH_URI,
                "token_uri": config.TOKEN_URI,
                "auth_provider_x509_cert_url": config.AUTH_PROVIDER_X509_CERT_URL,
                "client_x509_cert_url": config.CLIENT_X509_CERT_URL_WITHOUT_CLIENT_ID + client_id,
            }
        }
        flow = InstalledAppFlow.from_client_config(client_config, config.SCOPES)
        creds_obj = flow.run_local_server(port=0)
        sm.save_token(creds_obj, update_token=True)
        logger.info("新規認証が完了しました")
        return build(config.API_NAME, config.API_VERSION, credentials=creds_obj)
    except Exception as e:
        logger.error(f"認証に失敗しました: {e}")
        return None

AdMobのメディエーションレポートを取得

AdMobからは下記のデータを取得しています。

dimensions = [
   'DATE',                    # 日付
   'PLATFORM',                # プラットフォーム(iOS/Android)
   'FORMAT',                  # フォーマット(バナー/インタースティシャル)
   'AD_SOURCE',               # 広告ソース(AdMob Network等)
   'AD_UNIT',                 # 広告ユニット名
   'AD_SOURCE_INSTANCE',      # 広告ソースインスタンス名
   'MEDIATION_GROUP',         # メディエーショングループ名
]
metrics = [
   'ESTIMATED_EARNINGS',      # 推定収益(円)
   'AD_REQUESTS',             # 広告リクエスト数
   'CLICKS',                  # クリック数
   'IMPRESSIONS',             # インプレッション数
   'MATCH_RATE',              # マッチレート
   'IMPRESSION_CTR',          # インプレッションCTR
   'MATCHED_REQUESTS',        # マッチしたリクエスト数
   'OBSERVED_ECPM',           # 観測されたeCPM(円)
]

推定収益やeCPMは実際にAdMobの管理画面で表示される円換算された値をそのまま引っ張ってこれるので大変便利です。

AdMob APIの仕様上、一度に取得できる行が最大10万行に制限されているため取得期間を1年ごとに分割してリクエストを送信し、すべてのデータを取得できるようにしています。また、取得する収益データ(ESTIMATED_EARNINGS、OBSERVED_ECPM)は、精度を保つためにマイクロ単位(1/1,000,000)で提供されるため、/ 1000000で実際の金額に変換しています。これにより小数点以下の細かな収益も正確に記録できます。取得したデータはpandasのDataframeに格納されます。

mediation_report.py
import logging
import os
from datetime import datetime, timedelta
import pandas as pd
from google.cloud import bigquery as bq
import src.auth_manager as am
import src.config as config
import src.secret_manager as sm

logger = logging.getLogger(__name__)

def get_mediation_report(service, publisher_id, is_test):
    data_dir = config.DATA_DIR
    admob_mediation_report_file_name = config.ADMOB_MEDIATION_REPORT_FILE_NAME
    admob_mediation_report_file_path = os.path.join(data_dir, admob_mediation_report_file_name)
    # BigQueryからデータを取得
    client = bq.Client(config.PROJECT_ID)
    if is_test:
        table_id = f"{config.PROJECT_ID}.{config.SANDBOX_DATASET_ID}.test_admob_{config.BQ_TABLE_ID}"
    else:
        table_id = f"{config.PROJECT_ID}.{config.BQ_DATASET_ID}.{config.BQ_TABLE_ID}"
    try:
        client.get_table(table_id)
        logger.info(f"テーブルが存在します: {table_id}")
        try:
            # BigQueryから最新の日付を取得するクエリ
            query = f"""
            SELECT MAX(record_date) as last_date
            FROM `{table_id}`
            """
            query_result = client.query(query).result()
            last_date_row = list(query_result)[0]

            # BigQueryから取得した最新日付
            last_date = last_date_row.last_date
            if isinstance(last_date, str):
                last_date = datetime.strptime(last_date, '%Y-%m-%d').date()
            elif hasattr(last_date, 'date'):
                last_date = last_date.date()

            # 翌日から開始
            start_date = datetime.combine(last_date, datetime.min.time()) + timedelta(days=1)
            # UTCからJSTに変換(9時間追加)
            end_date = datetime.now() + timedelta(hours=9)

            # last_dateとend_dateが同じか、start_dateが未来の場合は、データが最新であることを表示して終了
            if last_date >= end_date.date():
                logger.info(f"データは最新です: {table_id}")
                return pd.DataFrame()  # 空のDataFrameを返す
            else:
                # 期間を1年ごとに分割
                date_ranges = []
                current_start = start_date
                while current_start < end_date:
                    current_end = min(
                        datetime(current_start.year + 1, current_start.month, current_start.day),
                        end_date
                    )
                    date_ranges.append({
                        'start_date': {
                            'year': current_start.year,
                            'month': current_start.month,
                            'day': current_start.day
                        },
                        'end_date': {
                            'year': current_end.year,
                            'month': current_end.month,
                            'day': current_end.day
                        }
                    })
                    current_start = current_end + timedelta(days=1)

        except Exception as e:
            return pd.DataFrame()  # エラー時は空のDataFrameを返す

    except Exception as e:
        logger.info(f"テーブルが存在しません: {table_id}")
        start_date = datetime(2020, 1, 1)
        end_date = datetime.now() - timedelta(days=2)
        # 期間を1年ごとに分割
        date_ranges = []
        current_start = start_date
        while current_start < end_date:
            current_end = min(
                datetime(current_start.year + 1, current_start.month, current_start.day),
                end_date
            )
            date_ranges.append({
                'start_date': {
                    'year': current_start.year,
                    'month': current_start.month,
                    'day': current_start.day
                },
                'end_date': {
                    'year': current_end.year,
                    'month': current_end.month,
                    'day': current_end.day
                }
            })
            current_start = current_end + timedelta(days=1)

    # メディエーション向けの詳細なレポート設定(実際に利用可能な値のみ)
    dimensions = [
        'DATE',  # 日付
        'PLATFORM',  # プラットフォーム(iOS/Android)
        'FORMAT',  # フォーマット(バナー/インタースティシャル)
        'AD_SOURCE',  # 広告ソース(AdMob Network等)
        'AD_UNIT',  # 広告ユニット名
        'AD_SOURCE_INSTANCE', # 広告ソースインスタンス名
        'MEDIATION_GROUP', # メディエーショングループ名
    ]

    metrics = [
        'ESTIMATED_EARNINGS',  # 推定収益(円)
        'AD_REQUESTS',  # 広告リクエスト数
        'CLICKS',  # クリック数
        'IMPRESSIONS',  # インプレッション数
        'MATCH_RATE',  # マッチレート
        'IMPRESSION_CTR',  # インプレッションCTR
        'MATCHED_REQUESTS',  # マッチしたリクエスト数
        'OBSERVED_ECPM',  # 観測されたeCPM(円)
    ]

    all_rows = []
    for date_range in date_ranges:
        # レポート仕様を作成
        report_spec = {
            'date_range': date_range,
            'dimensions': dimensions,
            'metrics': metrics,
            'sort_conditions': [
                {'dimension': 'DATE', 'order': 'DESCENDING'},
                {'metric': 'ESTIMATED_EARNINGS', 'order': 'DESCENDING'},
            ],
        }

        request = {'report_spec': report_spec}

        try:
            # メディエーションレポートを実行
            response = (
                service.accounts()
                .mediationReport()
                .generate(parent=f'accounts/{publisher_id}', body=request)
                .execute()
            )

            # レスポンスをDataFrameに変換
            if isinstance(response, list) and len(response) > 1:
                # データ行を抽出
                for item in response:
                    if 'row' in item:
                        row = item['row']

                        # ディメンション値を取得
                        date = row['dimensionValues']['DATE']['value']
                        platform = (
                            row['dimensionValues']
                            .get('PLATFORM', {})
                            .get('value', 'Unknown')
                        )
                        format = (
                            row['dimensionValues']
                            .get('FORMAT', {})
                            .get('value', 'Unknown')
                        )
                        ad_source = (
                            row['dimensionValues']
                            .get('AD_SOURCE', {})
                            .get('value', 'Unknown')
                        )
                        ad_unit = (
                            row['dimensionValues']
                            .get('AD_UNIT', {})
                            .get('value', 'Unknown')
                        )
                        ad_source_instance = (
                            row['dimensionValues']
                            .get('AD_SOURCE_INSTANCE', {})
                            .get('value', 'Unknown')
                        )
                        mediation_group = (
                            row['dimensionValues']
                            .get('MEDIATION_GROUP', {})
                            .get('value', 'Unknown')
                        )

                        # メトリック値を取得
                        earnings = (
                            float(row['metricValues']['ESTIMATED_EARNINGS']['microsValue'])
                            / 1000000
                        )  # マイクロ単位を円に変換
                        requests = (
                            int(row['metricValues']['AD_REQUESTS']['integerValue'])
                            if 'AD_REQUESTS' in row['metricValues']
                            else 0
                        )
                        clicks = (
                            int(row['metricValues']['CLICKS']['integerValue'])
                            if 'CLICKS' in row['metricValues']
                            else 0
                        )
                        impressions = (
                            int(row['metricValues']['IMPRESSIONS']['integerValue'])
                            if 'IMPRESSIONS' in row['metricValues']
                            else 0
                        )
                        match_rate = (
                            float(row['metricValues']['MATCH_RATE']['doubleValue'])
                            if 'MATCH_RATE' in row['metricValues']
                            else 0.0
                        )
                        impression_ctr = (
                            float(row['metricValues']['IMPRESSION_CTR']['doubleValue'])
                            if 'IMPRESSION_CTR' in row['metricValues']
                            else 0.0
                        )
                        matched_requests = (
                            int(row['metricValues']['MATCHED_REQUESTS']['integerValue'])
                            if 'MATCHED_REQUESTS' in row['metricValues']
                            else 0
                        )
                        observed_ecpm = (
                            float(row['metricValues']['OBSERVED_ECPM']['microsValue'])
                            / 1000000
                            if 'OBSERVED_ECPM' in row['metricValues']
                            else 0.0
                        )

                        all_rows.append(
                            [
                                date,
                                platform,
                                format,
                                ad_source,
                                ad_unit,
                                ad_source_instance,
                                mediation_group,
                                earnings,
                                requests,
                                clicks,
                                impressions,
                                match_rate,
                                impression_ctr,
                                matched_requests,
                                observed_ecpm,
                            ]
                        )

        except Exception as e:
            logger.error(f"期間 {date_range['start_date']} から {date_range['end_date']} のデータ取得中にエラーが発生しました: {str(e)}")
            continue

    if all_rows:
        # DataFrameを作成
        df = pd.DataFrame(
            all_rows,
            columns=[
                'date',
                'platform',
                'format',
                'ad_source',
                'ad_unit',
                'ad_source_instance',
                'mediation_group',
                'estimated_earnings',
                'ad_requests',
                'clicks',
                'impressions',
                'match_rate',
                'impression_ctr',
                'matched_requests',
                'observed_ecpm',
            ],
        )

        # 日付をdatetime型に変換
        df['date'] = pd.to_datetime(df['date'])
        return df
    else:
        return pd.DataFrame()

こちらで上記3つをまとめて実行

引数is_testがTrueならsandbox環境のテーブルを、Falseなら本番環境のテーブルを扱います。返り値は先ほど説明した通りDataFrameなので、このまま分析に利用することも、BigQueryにアップロードすることも自由にできます。

main.py
import pandas as pd
import mediation_report as mr
import auth_manager as am
import config
import secret_manager as sm

def get_mediation_report_df(is_test):
    credentials_dict = sm.get_secret("credentials")
    service = am.authenticate(credentials_dict)
    publisher_id = credentials_dict['ADMOB_PUBLISHER_ID']
    mediation_report_df = mr.get_mediation_report(service, publisher_id, is_test)
    return processed_report_df

if __name__ == "__main__":
    get_mediation_report_df(is_test=False)

Terraformでインフラ作成

実行するプログラムが完成したので、Terraformでサクッとインフラを作成します。スケジュールの設定はタイムゾーンがAsia/Tokyoで毎朝6時("0 6 * * *")にしました。このほかにgcr.tfやbigquery_dataset.tfも適宜設定してシステムの構築は完了です。また、下記(_locals.tf)のようにDockerイメージのパスをadmob_revenue_to_bigqueryとしてモジュール化することでterraform内のさまざまな場所で簡単に参照できるようになります。これにより設定の再利用性や保守性を高めています。
_locals.tf

locals {
  image = {
    admob_revenue_to_bigquery = "${var.default_region}-docker.pkg.dev/${var.project_id}/admob-revenue-to-bigquery/admob-revenue-to-bigquery:latest"
  }
}
cloudrun.tf
resource "google_cloud_run_v2_job" "admob_revenue_to_bigquery" {
  provider     = google-beta
  project      = var.project_id
  name         = "admob-revenue-to-bigquery"
  location     = local.default_region
  launch_stage = "BETA"

  template {
    labels = {
      "service" : "admob-revenue-to-bigquery"
    }
    template {
      containers {
        image = local.image.admob_revenue_to_bigquery
        env {
          name  = "ENV"
          value = var.env
        }
        resources {
          limits = {
            cpu    = "1000m"
            memory = "4Gi"
          }
        }
      }
      timeout         = "43200s"
      service_account = module.admob_revenue_to_bigquery.email
      max_retries     = 1
    }
  }

  lifecycle {
    ignore_changes = [
      launch_stage,
    ]
  }
}
secret_manager.tf
resource "google_secret_manager_secret" "admob_revenue_to_bigquery_credentials" {
  project   = var.project_id
  secret_id = "xxxxxxxxxxxxxxxxx"
  labels = {
    "app" : "admob_revenue_to_bigquery"
  }
  replication {
    auto {}
  }
}

resource "google_secret_manager_secret" "admob_revenue_to_bigquery_token" {
  project   = var.project_id
  secret_id = "xxxxxxxxxxxxxxxxx"

  labels = {
    "app" : "admob_revenue_to_bigquery"
  }

  replication {
    auto {}
  }
}
service_account.tf
module "admob_revenue_to_bigquery" {
  source      = "./modules/service_account"
  project_id  = var.project_id
  name        = "admob-revenue-to-bigquery"
  description = "AdMobの収益データをBigQueryに転送するサービスアカウント"
  roles = [
    "roles/bigquery.jobUser",
    "roles/bigquery.dataEditor",
    "roles/bigquery.readSessionUser",
    "roles/run.invoker",
    "roles/artifactregistry.reader",
    "roles/secretmanager.admin",
  ]
}
cloud_scheduler.tf
resource "google_cloud_scheduler_job" "admob_revenue_to_bigquery_scheduler" {
  project     = var.project_id
  name        = "admob-revenue-to-bigquery"
  region      = local.us_region
  description = "AdMobの収益をBigQueryに転送する"
  schedule    = "0 6 * * *"
  time_zone   = "Asia/Tokyo"
  http_target {
    http_method = "POST"
    // Google CloudのAPI経由でJobの実行を呼び出す
    uri = "https://${local.us_region}-run.googleapis.com/apis/run.googleapis.com/v1/namespaces/${var.project_id}/jobs/${google_cloud_run_v2_job.admob_revenue_to_bigquery.name}:run"
    oauth_token {
      service_account_email = module.admob_revenue_to_bigquery.email
    }
  }

まとめ

今回の記事では、AdMobの広告収益データをAPI経由で取得し、DataFrameに格納するまでのプロセスをご紹介しました。Cloud Runのスケジュール実行により、煩わしいスプレッドシートへの手作業転記はもう必要ありません。弊社ではこのデータをBigQueryに保存し、BIツールでダッシュボードを作成して収益状況を可視化することで、より良い議論ができるよう活用しています。また、機械学習モデルを構築し、過去のデータから将来の収益を予測してみるのも面白そうですね。
本記事が、皆さんのデータ活用の一助となれば幸いです。

WED Engineering Blog

Discussion