はじめに
こんにちは、株式会社ミツモアでデータエンジニアをしている酒井です。
ミツモアのデータチームでは、Digdag + dbt で運用していたデータパイプラインを Dagster + dbt へ移行しました。Dagster は日本語の記事自体がまだ少なく、特に運用レベルの実装パターンに踏み込んだものは見つかりにくかったため、サンプルの一つになればと思い本記事を書きました。
本記事では、ミツモアのデータチームが運用しているDagsterプロジェクトを題材に、
- ディレクトリをどう切っているか
Definitionsをどうマージしているか- dbt / リソース / スケジュールをどこにどう置いているか
- ローカル開発と本番運用の差分をどう吸収しているか
といった点を、実コードに寄せて紹介します。「Dagster の概念は分かったけど、自分たちのプロジェクトをスケールさせる書き方が見えてこない」という方の参考になれば幸いです。
なお、本記事では Dagster の基礎用語(asset / job / schedule / sensor / resource / definitions)は既知のものとして扱います。また「workflow」という語は Dagster の公式用語ではなく、本記事では 「1 つのデータ連携・処理のまとまり(= 1 つ以上の asset + それを束ねる job + トリガとなる schedule または sensor)」 を指す用語として使っています。
また、以降に登場するディレクトリ構成図やファイル名・コード例は、説明をわかりやすくするために実コードを一部簡略化・抽象化しており、本物のリポジトリと一字一句一致するわけではない点はあらかじめご了承ください。
対象とするプロジェクトの規模感と実行環境
参考までに、本記事で扱う Dagster プロジェクトの規模はおおむね次の通りです。
- workflowの数: 50 以上
- dbt modelの数: 1000以上
- 連携先データソース: MongoDB / PostgreSQL / 各広告 API(Yahoo / Criteo / LINE / Microsoft / SmartNews / Meta など)/ Stripe / SendGrid / Notion / Google Drive / その他 SaaS
- 本番実行環境: Kubernetes (EKS)

データオーケストレーションの概要。Dagsterで処理しているデータ・データの流れに絞っています。
1. プロジェクト全体のレイアウト
まず全体像から押さえます。Dagster プロジェクトは、データチームの単一モノレポの中に apps/transformation(dbt プロジェクト)と並べる形で apps/orchestration として配置しています。
<モノレポリポジトリ>/ ├── apps/ │ ├── orchestration/ ← この記事の対象 (Dagsterプロジェクト) │ └── transformation/ ← dbt プロジェクト │ ├── dbt_project.yml │ ├── profiles.yml # 本番環境とローカルで認証方法を切り替える仕組みを後述します │ ├── models/ │ └── ... └── ...
apps/orchestration/ 直下はこんな感じです。
apps/orchestration/
├── pyproject.toml # uv で依存管理。dbt-bigquery もここから入る
├── uv.lock
├── justfile
├── Dockerfile
├── entrypoint.sh # 本番 (K8s) のコンテナ起動スクリプト
├── dagster_home/ # ローカル用 DAGSTER_HOME (dagster.yaml を置く)
├── tests/
└── src/
└── orchestration/
├── definitions.py # Dagster の code location entrypoint
├── defs/ # 配下は (a) 1 workflow = 1 フォルダ と (b) 共通定義フォルダ の 2 種類
│ │
│ │ # --- (a) workflow 単位のフォルダ ---
│ ├── mongodb/ # sensor 駆動の例
│ │ ├── definitions.py
│ │ ├── assets.py
│ │ ├── config.py
│ │ ├── job.py
│ │ ├── sensor.py
│ │ ├── extractors/ # 統合固有のロジック (複数モジュールに分かれる場合)
│ │ └── README.md
│ ├── microsoft_ads/ # schedule 駆動の例
│ │ ├── definitions.py
│ │ ├── assets.py
│ │ ├── config.py
│ │ ├── <具体的な処理名>.py # 統合固有のロジック (単一ファイルで完結する場合)
│ │ ├── job.py
│ │ └── schedule.py
│ ├── proone/
│ ├── stripe/
│ ├── dbt_daily_06h00m/
│ ├── dbt_weekly/
│ │
│ │ # --- (b) 共通定義フォルダ (workflow ではない) ---
│ ├── dbt/ # dbt project と @dbt_assets の共通定義
│ ├── resources/ # 複数 workflow から参照しうる Dagster リソース
│ │ ├── definitions.py
│ │ ├── bigquery.py
│ │ ├── gcs.py
│ │ ├── mongo.py
│ │ ├── notion.py
│ │ ├── ... # 1 ファイル = 1 ConfigurableResource
│ │ ├── gcp_auth.py # GCP 認証情報のフォールバック
│ │ └── sync_cache.py # Dagster の PG を間借りした sync cache
│ └── ...
└── utilities/ # Dagster の語彙に依存しない共通ヘルパー
├── default_status.py
└── run_dbt.py
なお、Dagster には複数の code location を 1 つの UI / デプロイにまとめる workspace という仕組みもありますが、以下の理由から本プロジェクトでは採用していません。
- チームやデプロイサイクル・Python 依存を分ける必要がなく、複数 code location に切る動機がない
- 後述の
Definitions.merge()で「サブディレクトリごとに独立して書く」モジュール性は得られる
そのため src/orchestration/definitions.py を唯一の code location とし、その下で defs/<workflow名>/ ごとに Definitions を分割する方針に揃えています。
src/orchestration/ の中で役割を分けているのは次の 3 つです。
definitions.py: code location のエントリポイント。各サブディレクトリが export する Definitions をかき集めて 1 つにマージする。defs/: 本体。原則「defs/<workflow名>/1 つ = 1 workflow」という方針で分けています。例外としてdbt/(dbt 連携の共通定義)とresources/(複数 workflow で共有する Dagster リソース)の 2 つは、workflow ではなく共通定義の置き場です。utilities/:defs/をまたいで使い回す、Dagster の装飾子(@asset/@sensor等)に依存しない関数だけを置く場所。
dagster_home/ はローカル開発用です。本番 (Kubernetes) では entrypoint.sh が別途 DAGSTER_HOME をセットするため、リポジトリ内の dagster_home/ はローカル限定のセットアップ置き場として扱っています。
なお、以降のコード例で各スニペットの先頭にコメントで書いてあるファイルパスは、すべて apps/orchestration/src/orchestration/ からの相対パスです。(dbt profiles.yml だけは例外)
2. 各workflowフォルダ配下について
definitions.py — 各workflowフォルダのエントリポイント
同フォルダの assets.py / job.py / schedule.py / sensor.py を import して、最後に Definitions を 1 つ組み立てて definitions_<workflow> として export するだけのファイルです。
# defs/<workflow>/definitions.py from dagster import Definitions from orchestration.defs.<workflow>.assets import <workflow>_assets from orchestration.defs.<workflow>.job import <workflow>_job from orchestration.defs.<workflow>.schedule import <workflow>_schedule definitions_<workflow> = Definitions( assets=<workflow>_assets, # list[AssetsDefinition] jobs=[<workflow>_job], schedules=[<workflow>_schedule], )
ルートの src/orchestration/definitions.py がこれを Definitions.merge() で集めます。
config.py — 宣言的な設定
対象テーブル一覧・フィールドの命名マッピング・sync_type の選択といった「データだが定数として扱いたいもの」を dataclass / Enum / dict のリテラルで集めます。
例えば「対象テーブルのリスト」を持つ workflow の config.py はこんな雰囲気です。
# defs/<workflow>/config.py from dataclasses import dataclass @dataclass(frozen=True) class TableConfig: name: str sync_type: str # "full_replace" / "key_based" / "log_based" TABLES: list[TableConfig] = [ TableConfig("users", "key_based"), TableConfig("orders", "key_based"), TableConfig("events", "full_replace"), # ← 新規テーブルは行を 1 つ足すだけ ]
assets.py 側はこの TABLES をループして table ごとに 1 つの asset を動的生成します。
# defs/<workflow>/assets.py from dagster import asset, AssetsDefinition from orchestration.defs.<workflow>.config import TABLES def _build_asset(table: TableConfig) -> AssetsDefinition: @asset(name=table.name, group_name="<workflow>") def _asset(context): # 内部で table.sync_type を見て extractor を呼び分ける ... return _asset <workflow>_assets: list[AssetsDefinition] = [_build_asset(t) for t in TABLES]
設定とロジックを切っておくことで、対象テーブル 1 つ追加 = config.py に 1 行追加するだけ で済むようにしています。
extractors/ と <具体的な処理名>.py — 統合固有のロジック
assets.py から切り出した I/O とロジック本体の置き場です。
複数モジュールにまたがるなら extractors/ フォルダ、単一機能で完結するなら <具体的な処理名>.py をフォルダ直下に置きます。「最初は 1 ファイル直置きで始めて、増えてきたら extractors/ に昇格させる」ユルい運用です。
3. dbt 連携の置き方
dbt 連携は他の workflow と少し違って、「dbt project と asset の宣言」を共通フォルダに置き、「いつ・どのタグの dbt モデルを動かすか」を時刻バケットごとのフォルダに分けています。
defs/ ├── dbt/ # (b) 共通定義: DbtProject + @dbt_assets ├── dbt_bihourly/ # (a) workflow: tag:bihourly の dbt モデルを 2 時間おきに実行 ├── dbt_daily_06h00m/ # (a) workflow: tag:daily の dbt モデルを毎日 06:00 に実行 └── dbt_weekly/ # (a) workflow: tag:weekly の dbt モデルを週次で実行
defs/dbt/: dbt project と asset の共通定義
defs/dbt/definitions.py で DbtProject と @dbt_assets を一度だけ宣言します。各時刻バケットのフォルダはここから dbt_project_assets を import して、dbt_select で絞り込んで使う構造です。
# defs/dbt/definitions.py import os from pathlib import Path from dagster_dbt import DbtCliResource, DbtProject, dbt_assets from orchestration.defs.dbt.translator import CustomDagsterDbtTranslator from orchestration.defs.dbt.run_vars import DbtRunVars _is_production = os.environ.get("IS_PRODUCTION", "false") == "true" dbt_project = DbtProject( project_dir=Path(__file__).resolve().parents[5] / "transformation", target="production" if _is_production else "dev", ) dbt_project.prepare_if_dev() @dbt_assets( manifest=dbt_project.manifest_path, project=dbt_project, exclude="elementary", dagster_dbt_translator=CustomDagsterDbtTranslator(), ) def dbt_project_assets(context, dbt: DbtCliResource, config: DbtRunVars): # 実プロジェクトでは dbt run の vars 組み立てや log のハンドリングを # utilities/run_dbt.py にまとめている。ここでは省略。 yield from dbt.cli(["run", "--vars", "..."], context=context).stream()
dbt 実行 workflow フォルダ: dbt_bihourly, dbt_daily_06h00m, ...
中身は schedule + job の薄いフォルダです。ここでは dbt_bihourly を例にとります。
# defs/dbt_bihourly/job.py from dagster import define_asset_job from dagster_dbt import build_dbt_asset_selection from orchestration.defs.dbt.definitions import dbt_project_assets dbt_bihourly_job = define_asset_job( name="dbt_bihourly_job", selection=build_dbt_asset_selection( [dbt_project_assets], dbt_select="tag:bihourly", ), )
defs/dbt/ で宣言した dbt_project_assets から tag:bihourly の付いたモデルだけを取り出して 1 ジョブにまとめます。dbt 実行の直後に追加の Python asset を走らせたいケース(例: dbt mart を別 dataset にコピー)は、同フォルダに asset を足して AssetSelection.groups(...) | build_dbt_asset_selection(...) のように selection を union するだけです。
4. ローカル / 本番の挙動差分
50 超のworkflowを抱えるプロジェクトをローカルで just serve 起動するときに困るのが、ローカルでも schedule / sensor が動いてしまう ことです。MongoDB / Slack / 各種 SaaS の本番資格情報をローカルに置いていない(置きたくない)ので、自動起動したら大量の失敗が出ます。
以下はこのような「ローカル / 本番の挙動差分」を吸収するための実装です。
Dagster の default_status ヘルパー
utilities/default_status.py に共通ヘルパーがあります。
# utilities/default_status.py import os from dagster import DefaultScheduleStatus, DefaultSensorStatus def _is_local() -> bool: # K8s pod でも CI runner でもなければローカル扱い return not os.environ.get("KUBERNETES_SERVICE_HOST") and not os.environ.get("CI") def default_schedule_status() -> DefaultScheduleStatus: return DefaultScheduleStatus.STOPPED if _is_local() else DefaultScheduleStatus.RUNNING def default_sensor_status() -> DefaultSensorStatus: return DefaultSensorStatus.STOPPED if _is_local() else DefaultSensorStatus.RUNNING
判定の鍵は環境変数 2 つだけです:
KUBERNETES_SERVICE_HOST: K8s pod 内でセットされるCI: GitHub Actions がセットする
どちらも立っていなければローカル → schedule / sensor は STOPPED で初期化。K8s pod や CI 上では RUNNING になります。
dbt の profiles.yml: OAuth ↔ service-account 自動切替
dbt 側も同じ「ローカル ↔ K8s / CI」の切り分けが必要です。apps/transformation/profiles.yml の各 target で method を Jinja 式で動的に決めて います。
# apps/transformation/profiles.yml default: outputs: dev: type: bigquery method: "{{ 'oauth' if not env_var('KUBERNETES_SERVICE_HOST', '') and not env_var('CI', '') else 'service-account' }}" project: example_project_id dataset: example_dataset_name keyfile: "{{ env_var('GCP_CREDENTIALS_FILE', '/opt/dagster/transformation/gcp-credentials.json') }}"
- ローカル:
method=oauth→gcloud auth application-default login済みの個人 Google アカウントで BigQuery を叩く。keyfile行は dbt-bigquery 側で無視される - K8s pod / CI:
method=service-account→keyfileで指定された SA keyfile を読む
これで「ローカルからの実行は BigQuery audit log に個人アカウントが残る(事故時の追跡可能性)」「本番デプロイは SA 認証で個人に依存しない」を実現しています。
5. まとめ
本記事ではミツモアのデータチームで運用している Dagster プロジェクトのおおまかな構成を紹介しました。 本記事の内容がそのまま他のプロジェクトに当てはまるとは限りませんが、Dagster でデータ基盤を組むときのサンプルの一つとして参考になれば幸いです。
採用情報
ミツモアでは、データやAIを活用してデータドリブンな風土のある会社にて一緒に働く仲間を募集中です。
「技術で課題を解くことにワクワクできる人」や「仕組みで社会を良くしたい、そんなエンジニアになりたい人」、ぜひご応募をお待ちしています!
ミツモア採用ページ: https://corp.meetsmore.com/