コンテンツにスキップ

Databricks 連携

アンケートの回答データは Postgres ではなく Databricks の Unity Catalog にある。

経緯

以前は Creative Survey が S3 に置く answers.csv.gz / panels.csv.gz を arrange のバッチが Postgres にキャッシュしていた。しかし同じバケットの同じファイルを Databricks 側のパイプラインも取り込んでいたため、同一データを二重に管理していた。

さらに Databricks 側は回答とパネルを結合し調査ごとに分割するところまで加工済みだったため、arrange はそれを読むだけでよくなった。

廃止したもの 経緯
s3_answers / s3_panels / s3_sync_logs 回答本体のキャッシュが不要になった
batch/s3-sync.ts の syncMeta / syncAnswers 同期そのものが不要になった
lib/s3-client.ts、S3_* 環境変数 S3 を読まなくなった
/api/v1/s3/sync-* 系 5 ルート 同期の起動・状況確認が不要になった

テーブル

テーブル 用途
cs.cs_dm.dm_answers_<survey_id> 調査ごとの回答。回答とパネルが LEFT JOIN 済みなので、こちら側で結合する必要がない。611 調査ぶん存在する
cs.cs_dwh.dwh_answers / dwh_panels 全調査横断。調査メタの集計にだけ使う

cs_dm は PII を含む。cs_dm_pub は氏名・メール・電話番号などを *** にマスクした版だが、候補者への連絡にそれらが必要なので arrange は cs_dm を使う。

生成パイプライン

Databricks 側の cs_data_pipeline_s3_delivery(毎日 5:00:45 JST、intage 管理)が生成する。

cs_download_files_from_cs_s3     S3 から DL
  ├── cs_load_answers            → dl_answers
  └── cs_load_panels             → dl_panels
        └── cs_split_survey_tables      → cs_dm.dm_answers_<survey_id>
              └── cs_create_non_pii_tables → cs_dm_pub.dm_pub_answers_<survey_id>

調査別テーブルは差分更新

cs_split_survey_tables は直近 3 日に更新のあった調査だけを作り直す(UPDATED_SURVEY_ID_DAYS = 3)。締め切りから 3 日以上経った調査のテーブルは更新されない。

実装

ファイル 役割
src/lib/databricks-client.ts SQL Statement Execution API を fetch だけで叩くクライアント
src/repositories/databricks.ts 上記テーブルへのクエリ
src/batch/databricks-meta-sync.ts databricks_surveys の日次更新

クライアントの設計

  • ドライバ依存を持たない。 @databricks/sql は Thrift 依存で Bun での動作リスクがあるため、REST API を fetch で直接叩く
  • 結果は常に EXTERNAL_LINKS チャンクで受ける。 INLINE は 25 MiB 上限で、1 調査が 50 万行を超えるため収まらない。チャンクの presigned URL には認証情報が埋まっているので Authorization ヘッダを付けてはいけない
  • getEnvConfig() を通さず process.env を直読みする。 バッチはアプリ用の必須環境変数が揃わない環境で動くため(lib/s3-client.ts が S3_URI を直読みしていたのと同じ理由)

値の型

JSON_ARRAY 形式で受けるためすべて文字列で返る。

列 返る値
is_answered / is_completed / is_mobile '1' / '0'
timestamp 列 ISO 8601 文字列

repositories/databricks.ts の parseFlag / parseTimestamp / parseCount で変換する。

survey_id の扱い

テーブル名に埋め込むためパラメータバインドできない。数値のみを許可する検証(assertSurveyId)を必ず通す。

性能

操作 実測
1 クエリの固定オーバーヘッド 約 1 秒
調査メタの全件集計(68M 行) 約 15 秒
単一調査の読み出し(531,221 行) 約 5 秒
全行取得(531,221 行、複数チャンク) 約 165 秒

固定オーバーヘッドが大きいため、バッチ処理は 1 クエリで多く取る。survey-import.ts の BATCH_SIZE は 1000。20,712 パネルの調査での実測は次のとおり。

バッチサイズ 1 バッチ バッチ数 全体
50 1.53s 415 10.6 分
500 4.24s 42 3.0 分
1000 6.60s 21 2.3 分
2000 13.04s 11 2.4 分

認証

OAuth M2M(サービスプリンシパルの client credentials)。トークンは有効期限までキャッシュする。

DATABRICKS_HOST
DATABRICKS_WAREHOUSE_ID
DATABRICKS_CLIENT_ID
DATABRICKS_CLIENT_SECRET
DATABRICKS_CATALOG        既定 cs
DATABRICKS_SCHEMA_DM      既定 cs_dm
DATABRICKS_SCHEMA_DWH     既定 cs_dwh

ローカルでは CLI のトークンで代替できる。

DATABRICKS_TOKEN=$(databricks auth token -p <profile> | jq -r .access_token)

DATABRICKS_TOKEN があればそれを使い、無ければ OAuth M2M でトークンを取得する。

サービスプリンシパルには Unity Catalog の権限が必要

cs カタログに USE_CATALOG / USE_SCHEMA / SELECT が必要。本来は accs_uc_cs_cs_dm_readonly と accs_uc_cs_cs_dwh_readonly のグループに追加するのが筋だが、グループのメンバー変更はアカウントレベルの操作になる。

データ構造で注意すべき点

設問の識別子は設問文そのもの

mart に設問 ID が無い。フィールドは 設問文、マトリクス設問は 設問文|||選択肢文 の複合キーで識別する。

同じ設問文が複数の親質問の下に現れる

「その理由を教えてください。」のような自由記述は複数の設問に付随するため、キーが衝突する。sub_item_sentence も空で、区別する手段が無い(区別できるのは行 ID と時刻だけ)。

そのため回答を畳むときは、1 パネルが同一フィールドに 2 件以上の回答を持つフィールドを実データから検出し、そのフィールドだけ配列として保持する(fetchMultiValueFields)。

answer_type からは単複を判別できない

最多の型は テキスト選択(単・複)(4,242 万行 / 8,851 設問)で、文字どおり「単一または複数」。型名で判定しようとすると必ず破綻する。

複数選択は「選択肢 1 つ = 1 行」

複数選択の設問は、選択された選択肢ごとに 1 行が入り、answer_item_sentence と value の両方に選択肢名が入る。つまり各選択肢が独立したフィールドになる。

一方プルダウンやマトリクスでは answer_item_sentence がサブ設問名で value が回答値になる。answer_item_sentence === value かどうかで振り分けられる。

answer_type の一覧(実測)

answer_type 行数 設問数
テキスト選択(単・複) 42,423,070 8,851
テキスト入力 8,754,475 4,026
プルダウン 8,310,220 664
テキストマトリクス 6,310,701 661
画像選択(単・複) 1,310,490 578
スライドバー 953,629 421
NPSタイプ 76,835 102
サバイバルテキスト 11,867 9
固定回答タイプ 8,355 1
カレンダー選択 895 4
サバイバルイメージ 793 2

バッチの実行

Dockerfile.sync の ENTRYPOINT は ["bun", "run"] までで、スクリプトのパスと JSON 引数は ECS の container override で渡す。1 つのイメージで 2 つのバッチを起動できる。

# 調査メタの日次更新
aws ecs run-task \
  --cluster dev-heineken-interview-arrange-backend-sync \
  --task-definition dev-heineken-interview-arrange-backend-sync \
  --launch-type FARGATE \
  --network-configuration "awsvpcConfiguration={subnets=[<private-subnet>],securityGroups=[<lambda-sg>],assignPublicIp=DISABLED}" \
  --overrides '{"containerOverrides":[{"name":"sync","command":["src/batch/databricks-meta-sync.ts","{}"]}]}' \
  --region ap-northeast-1

ローカルでは直接実行できる。

bun run src/batch/databricks-meta-sync.ts '{"dryRun":true}'

定期実行はまだ設定されていない

cs_data_pipeline_s3_delivery の後に上記の ECS タスクを起動する想定だが、その連携は未設定。手動で実行しない限り databricks_surveys は更新されない。