🎻

dbt Platform (dbt Cloud)とオーケストレーションツールを組み合わせる

に公開

この記事はdbtアドベントカレンダー2025に寄稿しています。

はじめに

アナリティクスエンジニアのチームでは、基本的にdbtモデルの作成などのデータモデリングや分析支援がメインの業務になりますが、時折、

  • プロダクトで使っているCloud SQLのようなデータベースにデータを送ってほしい
  • 機械学習モデルにデータを渡して予測を行い、その予測データを含めて集計したい

などの、SQLだけでは書けない処理を実行したいというニーズが発生します。

Cloud Run Jobsなどで専用のETLバッチ処理を書いて実行する、といった方法も当然あるのですが、必ずしもアナリティクスエンジニアは、全員がソフトウェアエンジニアのバックグラウンドを持っている訳ではありません。Dockerコンテナのような新しいツールを導入して、カスタマイズされたETLバッチ処理をメンテナンスしていくのは、チームのスキルセットとして難しい場合があります。

そのため、普段から使っている既存のツールを組み合わせて、メンテナンスするコードを最小限にした形で実現するのが理想的です。

今回は、以下のようなチーム構成を仮定し、チーム間で連携しながらdbt Platformとオーケストレーションツールを組み合わせて使う方法について解説します。

アナリティクスエンジニアのチーム

  • 普段は dbt Platform をメインで使っている

データエンジニアのチーム

  • 普段はAirflowやPrefect、Dagsterのようなオーケストレーションツールをメインで使っていて、様々なETLバッチ処理を運用している


協力してみましたのイメージ by Gemini

dbt Platform の制約

dbt Platform (dbt Cloud)はご存知の通り、dbtのマネージドサービスとして、データ変換処理の構築に関して高度な機能を数多く提供しています。

dbt Platform上ではjobというリソースで複数のdbtコマンドをまとめて実行することができますが、dbtはあくまでデータ変換に特化したツールであり、それ以外の機能は限定的な側面があります。

例えば、job同士は連結が可能なため、簡単な依存関係を組むことはできますが、「jobに定義された日付変数をループで変化させながら実行する」といったオーケストレーションツールが得意な処理はできません。

こうしたdbtやdbt Platform単体では対応できないようなユースケースでは、dbt側もオーケストレーションツールと組み合わせて使う方法を公式ドキュメント上で紹介しています。

オーケストレーションツール側からdbt Platformのjobを実行する

多くのオーケストレーションツールでは、dbt Coreとdbt Platformの両方に対応したOperatorやTaskが提供されています。

今回はアナリティクスエンジニアチームがdbt Platformを使っているという想定なので、そのAPIを経由してdbt Platform上にあるjobをオーケストレーションツール側から実行する方法を採用します。

dbt Platformを使用することで、dbtで実行するコマンド等の情報はすべてdbt Platform側に定義されるため、オーケストレーションツール側のソースにdbtのコードを含めなくて良いなど、よりシンプルな形で実装することができます。

また、オーケストレーションツールとしては代表的なAirflowを例に解説します。Airflowでは、apache-airflow-providers-dbt-cloud というパッケージをAirflowに追加する必要があります。

ちなみに、Google CloudのCloud Composerであれば、既にこのパッケージが元からインストールされているため、追加のインストール作業自体も不要です。

https://docs.cloud.google.com/composer/docs/composer-versions

Airflow と dbt Platform の接続設定

初期設定として、Airflow側でdbt Platformとの接続設定を行います。

dbt PlatformのAPIを利用するためには、Service Tokenと呼ばれる認証情報を発行する必要があります。

https://docs.getdbt.com/docs/dbt-cloud-apis/service-tokens

今回はJobのトリガーにのみ使用するため、権限としては、Job RunnerJob Admin が付与されていれば問題ありません。

https://airflow.apache.org/docs/apache-airflow-providers-dbt-cloud/stable/connections.html

認証情報以外の設定値としては、Host URLやAccount IDなども必要になります。

これらすべてを dbt_cloud_default という名前のConnection IDとしてAirflowの接続設定にWeb UIなどから登録します。

Airflow DAGの実装例

接続設定が完了したので実際にパイプラインを構築します。今回は、機械学習の予測処理を含むパイプラインを例にして、

  • 機械学習の予測バッチに必要なデータを作成するためのjobをdbt Platformで実行
  • Vertex AIのBatch Predictionジョブを実行
  • 予測結果とともに、さらに追加の集計を行うjobをdbt Platformで実行

という3つのステップからなるパイプラインを構築します。

なお、以下のコードはあくまで最小限の実装になっており、実際のコードでは複数のパイプラインを定義する場合などに備えて、モジュール化などの処理を加えた上での実装を行うのが望ましいです。

from datetime import datetime
from airflow import DAG
from airflow.providers.dbt.cloud.operators.dbt import DbtCloudRunJobOperator
from airflow.providers.google.cloud.operators.vertex_ai.batch_prediction_job import CreateBatchPredictionJobOperator

# 設定値
DBT_CLOUD_CONN_ID = "dbt_cloud_default"
GCP_PROJECT_ID = "your-gcp-project-id"
GCP_LOCATION = "us-central1"
VERTEX_MODEL_NAME = "projects/your-gcp-project-id/locations/us-central1/models/your-model-id"

# dbt Job IDs
DBT_JOB_ID_PREPROCESS = 12345
DBT_JOB_ID_POSTPROCESS = 67890

with DAG(
    dag_id="dbt_vertex_pipeline_minimal",
    start_date=datetime(2024, 1, 1),
    schedule_interval=None,
    catchup=False,
    tags=["dbt", "vertex_ai"],
) as dag:

    # 1. dbt Cloud Job (前処理) を実行
    dbt_job_preprocess = DbtCloudRunJobOperator(
        task_id="dbt_job_preprocess",
        dbt_cloud_conn_id=DBT_CLOUD_CONN_ID,
        job_id=DBT_JOB_ID_PREPROCESS,
    )

    # 2. Vertex AI で Batch Prediction を実行
    # DWH が BigQuery なので、入出力に BigQuery を指定
    vertex_batch_predict = CreateBatchPredictionJobOperator(
        task_id="vertex_batch_predict",
        project_id=GCP_PROJECT_ID,
        location=GCP_LOCATION,
        job_display_name="dbt-vertex-prediction-job",
        model_name=VERTEX_MODEL_NAME,
        # 入力データ (BigQuery)
        instances_format="bigquery",
        bigquery_source="bq://your-project.dataset.input_table",
        # 出力先 (BigQuery)
        predictions_format="bigquery",
        bigquery_destination_prefix="bq://your-project.dataset.output_table",
    )

    # 3. dbt Cloud Job (後処理) を実行
    dbt_job_postprocess = DbtCloudRunJobOperator(
        task_id="dbt_job_postprocess",
        dbt_cloud_conn_id=DBT_CLOUD_CONN_ID,
        job_id=DBT_JOB_ID_POSTPROCESS,
    )

    # 依存関係の定義
    dbt_job_preprocess >> vertex_batch_predict >> dbt_job_postprocess

Airflowで書くと、このようにかなり少ないコード量で実装することができますが、すべて自力で実装しようとするとかなりのコード量になってしまいます。

特にdbt Platformは公式sdkライブラリがありませんので、直接requestsライブラリなどを使ってAPIを叩く必要があり、jobのステータス監視やリトライ処理なども自分で実装することになります。

その場合、ビジネスロジックと関係ない部分のコードが膨れ上がってしまい、バグのリスクが増大したりと、本質的な価値提供に集中する時間を奪われてしまいます。

ですが、Airflowのようなオーケストレーションツールには、こうした汎用的な機能も組み込まれた専用のOperatorが用意されているため、その点もツールを使うメリットの1つになります。

Operator内部の実装を見てみると、まず、self.hook.trigger_job_run メソッドでdbt Platformのjobを実行しますが、その後に self.hook.wait_for_job_run_status メソッドでjobがSuccessになるまで待機する処理が実装されています。

https://github.com/apache/airflow/blob/providers-dbt-cloud/4.5.0/providers/dbt/cloud/src/airflow/providers/dbt/cloud/operators/dbt.py#L182-L215

self.hook.wait_for_job_run_status メソッドの中身を見てみると、jobのステータスを定期的にポーリングして、期待値と合致する状態(この場合はSuccess)になるまで待機する処理になっていることが分かります。

https://github.com/apache/airflow/blob/providers-dbt-cloud/4.5.0/providers/dbt/cloud/src/airflow/providers/dbt/cloud/hooks/dbt.py#L781-L824

また、もしdbt Platform側のAPIのバージョンが上がった場合でも、Operatorを使用していれば、パッケージのバージョンアップを行うだけで最新のAPIに対応できるため、メンテナンス性の面でも優れています。

チーム間連携の容易さ

また、この構成はコードのメンテナンス性が高いだけでなく、チーム間連携の面でも優れています。

チーム間の連携は一歩間違えるとボトルネックになってしまうため、以下の観点で設計することが重要です。

  • 新しい技術スタックに触れるアナリティクスエンジニアの学習コストを最小限にする
  • 他のチームにリソースを共有するデータエンジニアチームの運用・サポートに関する負担をなるべく減らす

デプロイ承認プロセスの簡略化

今回のようなビジネスロジックが多分に含まれるパイプラインでは、一度デプロイしたら終わりというものではなく、往々にして継続的に改修が発生します。例えば、機械学習モデルが更新されて新しい特徴量が必要になるといったケースです。

もしdbtのコードがAirflow側に含まれている場合(dbt CoreをAirflow上で直接実行する構成の場合)、dbtのコードを修正した際にはAirflow側のリポジトリに対してデプロイを行う必要が出てきます。その場合、改修を行ったアナリティクスエンジニアは、自分のチームメンバーに加えて、Airflowを管理しているデータエンジニアチームにもレビューを依頼する必要があります。データエンジニアとしても要件が良く分かっていないビジネスロジックのコードをレビューするのは負担が大きくなります。

一方、dbt Platform側にコードを配置している今回の構成では、dbt側のコードを修正してdbt Platform上にデプロイするだけで完結します。そのため、Airflowを管理しているデータエンジニアチームに対して、毎回レビューを依頼したり、PRを送る必要がありません。

ローカル開発環境構築を不要にする

前述の通り、アナリティクスエンジニアはAirflowのコードを触る必要がないため、ローカルにAirflowの開発環境を構築する必要がなくなります。

どの環境でAirflowを動かしているかにもよりますが、例えばCloud Composerを使っている場合、composer-local-devというパッケージは用意されているものの、それでもセットアップや動作確認の作業はそれなりに複雑で時間もかかります。

Dockerを使ったことがない、セットアップから始めないといけないような場合、アナリティクスエンジニアチームにとってはハードルが高くなってしまいます。

そのため、アナリティクスエンジニアはAirflowの学習をするとしてもWeb UIの使い方や基本的な概念の理解に集中した方が、双方のチーム全体としての生産性が向上するでしょう。

コスト

良い点ばかり話してきましたが、もちろんこのアーキテクチャにもデメリットはあります。その1つがインフラコストです。

今回のケースは既に導入済のツールを組み合わせるという形を取っているため、追加のコストは発生しませんが、新規に導入する場合は、dbt Platformの利用料金に加えて、オーケストレーションツールのインフラコストも発生します。

どちらもエンタープライズ向けの機能を備えたツールであり、コストもそれなりにかかります。そのため、両方のツールを1つの小規模なチームが契約・運用するのはコスト的に厳しい場合があります。

ETL処理を担当しているデータエンジニアだけで社内データチームが構成されており、アナリティクスエンジニアはまだ採用していないような場合は、dbt Platformを新規に導入するよりも、Astronomer社が開発しているCosmosのようなライブラリや、Dagsterのようなdbt Coreとの親和性が高いツールを活用する方がコスト観点では適しているでしょう。

さいごに

今回は、dbt Platform とオーケストレーションツールを組み合わせて使う方法について解説しました。

コード自体は今やAIがあっという間に書いてくれるようになり、どんどん精度も上がってきています。ただ、そのコードはどうすればメンテナンスしやすいか、何を実行基盤とするのが適切か、といった設計の部分は、まだ人間の腕の見せ所なのかもしれません。

Discussion