Dagster Orchestration
Dagster is the execution runtime for production training runs. The module-level defs object in src/daedalus/definitions.py is the entrypoint that Dagit, the dg CLI, and dagster dev load:
# src/daedalus/definitions.py
defs = dg.Definitions(
assets=[*feature_view_specs, pixai_feed_relevance_training],
jobs=[training_job],
schedules=[training_schedule],
resources={"cfg": DaedalusConfig()},
)Everything below is defined in src/daedalus/defs/training.py.
The hybrid model: a partitioned graph-backed asset
The top abstraction is a daily-partitioned asset per training service — one partition per day, which is also the schedule cadence. The partition is produced by a graph of ops generated from the compiled operator DAG's engine units (pipeline.units()) — so the graph mirrors the pipeline, not a fixed two-stage split. One op per engine unit, visible in Dagit and individually retriable:
compile_pipeline ─▶ source ─▶ enrich ─▶ (asset materialization)The ops are factory-generated per service from the unit list, so a service with no pythonic unit (no declared embedding columns) yields just compile_pipeline ─▶ source — no enrich op. Each op wraps the executor (pipelines.executor.run_pipeline), selecting its engine unit via only_engines — no feature logic is re-derived. Each emits rich metadata so a run is self-describing in Dagit:
| Op | Engine unit | Metadata emitted |
|---|---|---|
compile_pipeline | — | operators count, the full dag as JSON |
source | sql / sql_arrow_udf (only_engines=SOURCE_ENGINES) | rows, output path |
enrich | pythonic (only_engines=ENRICH_ENGINES) | rows, columns, embedding_columns, output_path, partition |
The enrich op reads counts and schema from the parquet metadata only — it never materializes the large 1152-d embedding partition into memory just to report stats.
Feature views as external assets (lineage)
Every feature view in the registry is declared as a Dagster external asset, so the asset graph is the lineage. They appear in the feature_views group keyed as feature_view/<name> with the parquet kind:
dg.AssetSpec(
key=dg.AssetKey(["feature_view", name]),
group_name="feature_views",
kinds={"parquet"},
)Job + schedule
training_job = dg.define_asset_job(
name="pixai_feed_relevance_daily",
selection=[pixai_feed_relevance_training],
)
training_schedule = dg.build_schedule_from_partitioned_job(training_job)build_schedule_from_partitioned_job derives a daily schedule directly from the daily partition definition, so production runs follow the partition cadence with no separate cron string to maintain.
The DaedalusConfig resource
Where the engine reads config and writes output is a ConfigurableResource, overridable per run or per environment:
class DaedalusConfig(dg.ConfigurableResource):
runtime_config_path: str = "feature_services/pixai_feed_relevance/settings.yaml"
aggregation_config_path: str = "config/training/aggregations.yaml"
feature_views_dir: str = "feature_views"
feature_services_dir: str = "feature_services"
output_root: str = "data/training_output"
enriched_output_root: str = "data/training_output_enriched"
# Local Ray envelope (prod overrides to 16c / 64 GiB).
ray_num_cpus: int = 16
ray_memory_gb: int = 44
ray_object_store_gb: int = 12Daily today, hourly later
The cadence is daily end to end. Moving to hourly is a one-line swap to HourlyPartitionsDefinition once the feed supports it.
Running it
# Launch Dagit (the web UI) against the definitions module
uv run dagster dev -m daedalus.definitions
# Validate the definitions load cleanly
uv run dg check defsAgent CLI: ad-hoc materialization & lineage
For agents and quick checks, two top-level daeda commands wrap the same Dagster asset in-process (src/daedalus/commands/agent.py) — no running Dagit instance required.
daeda materialize-day
Materializes one training-day partition (via dg.materialize on pixai_feed_relevance_training) so you can inspect sample output. Only the pixai_feed_relevance service is supported.
uv run daeda materialize-day pixai_feed_relevance 2026-02-01It prints a JSON summary pulled from the materialization metadata:
{
"service": "pixai_feed_relevance",
"date": "2026-02-01",
"rows": 1234567,
"output_path": "data/.../dt=2026-02-01",
"embedding_columns": ["image_embedding", "like_artwork_avg_embeds"]
}daeda lineage
Shows upstream/downstream lineage for one feature view or service, resolved from the platform registry (not Dagster):
uv run daeda lineage pixai_feed_relevance # table (default)
uv run daeda lineage user_likes --format jsonFor a service it lists upstream views and their sources; for a view it lists upstream sources and downstream services.
See also
- Training Pipeline overview — what the ops actually run.
- JSON-RPC API — the
runs.*andassets.lineagemethods that reach the same Dagster job over HTTP. - CLI reference.