Dagster operations
In plain English: AlphaSwarm’s production Dagster graph is a classic
Definitions object, not a dg / create-dagster project. You load
one of two documented live code locations depending on where you
run: compose viz loads the monolith module; the Kubernetes user-code
overlay loads the platform pipelines module. They are not a dual-load
of the same file. Instance limits live in dagster.yaml. The
interactive try-it sandbox is a separate page.
Entrypoint and validate gate
Pin is dagster==1.13.13 (pyproject.toml extra dagster). dg is
not installed or scaffolded. The working path is:
python -m dagster definitions validate -m alphaswarm.dagster.definitions
The monolith entrypoint is
alphaswarm/dagster/definitions.py
(defs = Definitions(...)). Settings default
dagster_module_path is alphaswarm.dagster.definitions.
Two live code locations
These are two Definitions graphs. Activating both in one Dagster
instance is a later project (not wired today).
| Location | Module | How it is launched | What it owns |
|---|---|---|---|
| Monolith | alphaswarm.dagster.definitions | Compose viz: dagster dev -h 0.0.0.0 -p 3001 -m alphaswarm.dagster.definitions in compose/docker-compose.viz.yml. Settings: dagster_module_path. | Ingest, entities, catalog, Airbyte health, platform ops, Alpha Vantage intraday, freshness, run-failure |
| Platform | pipelines.dagster_user_code.definitions | Helm overlay values-pipelines-user-code.yaml: dagsterApiGrpcArgs: ["-m", "pipelines.dagster_user_code.definitions"] on deployment pipelines-user-code (image ghcr.io/julianwiley/alphaswarm-pipelines). Compose viz does not load this module. | MinIO/CDC/vectorize/DataHub/RAG/Alpha Vantage assets and their jobs/schedules |
The documented live platform location is the Helm overlay
values-pipelines-user-code.yaml (pipelines-user-code →
pipelines.dagster_user_code.definitions). Default Helm
values.yaml still ships a third user-code deployment
bootstrap-user-code with --python-file /opt/dagster/user_code/definitions.py; that file is neither graph,
so installing values.yaml without the overlay does not load
pipelines.dagster_user_code.definitions.
The monolith module docstring still says Helm pipelines-user-code
loads dagster api grpc -m alphaswarm.dagster.definitions. That
command is not what the Helm overlay runs. Follow the table above.
Platform jobs include minio_to_postgres_job, vectorization_job,
cdc_job, hybrid_transform_job, pdf_ingest_job, csv_ingest_job,
graphrag_job, the DataHub job family, and
alphaswarm_alphavantage_intraday_*. Schedules:
cdc_hourly_schedule, vector_daily_schedule,
datahub_daily_schedule,
alphaswarm_alphavantage_intraday_delta_schedule. Source:
pipelines/dagster_user_code/definitions.py.
Monolith inventory (live Definitions)
Counts from ALL_JOBS / ALL_SCHEDULES / ALL_SENSORS /
ALL_ASSET_CHECKS / FRESHNESS_CHECKS / all_assets() after the
2026-08-13 expand/enhance/harden work. airbyte_connections and dbt
groups stay empty unless an Airbyte workspace or dbt mesh is
configured. partitions.py
is unused.
| Kind | Count | Names |
|---|---|---|
| Jobs | 12 | full_data_refresh_job, regulatory_refresh_job, entity_extraction_job, compaction_job, profiling_job, datahub_sync_job, pipeline_manifest_materialization_job, platform_ops_job, time_partitioned_sources_job, alphavantage_intraday_partition_job, materialize_features, alphavantage_intraday_delta_job |
| Schedules | 10 | daily_full_refresh 0 2 * * *; weekday_regulatory_refresh 0 4 * * 1-5; hourly_datahub_sync 15 * * * *; six_hourly_profiling 30 */6 * * *; weekly_compaction 0 5 * * 0; daily_entity_enrichment 0 6 * * *; daily_time_partitioned_sources 10 3 * * *; daily_alphavantage_intraday_partition 40 1 * * *; daily_platform_ops 0 1 * * *; alphavantage_intraday_delta_schedule 20 * * * * |
| Sensors | 3 | pipeline_manifests_changed, alphaswarm_run_failure, alphaswarm_automation_condition_sensor |
| Asset checks | 6 | datahub_platform_instance_is_aqp, datahub_external_platforms_exclude_assistants, iceberg_namespace_configured, iceberg_medallion_namespace_prefix_matches_layer, iceberg_compaction_non_empty, airbyte_bronze_landing_namespace |
| Freshness checks | 2 | last-update checks grouped by tightest ingest-cron interval |
| Assets | 36 | alphaswarm_sources 9, alphaswarm_entities 6, alphaswarm_catalog 3, alphaswarm_profiling 1, alphaswarm_compaction 1, airbyte 2, alphaswarm_engine 1, alphaswarm_platform_ops 6, platform 4, alphaswarm_alpha_vantage 3 |
dagster.yaml
Source:
alphaswarm/dagster/dagster.yaml.
Schema-valid for 1.13.13 — no named concurrency.pools.config
block.
| Field | Value |
|---|---|
run_retries.max_retries | 3 |
run_monitoring.free_slots_after_run_end_seconds | 300 |
concurrency.pools.default_limit | 8 |
concurrency.runs.max_concurrent_runs | 8 |
Compose. Bind-mount
../../alphaswarm/alphaswarm/dagster/dagster.yaml →
/app/data/dagster_home/dagster.yaml:ro.
DAGSTER_HOME=/app/data/dagster_home.
Kubernetes. Helm
values.yaml
sets global.dagsterHome: /opt/dagster/dagster_home. The chart renders
the instance ConfigMap and mounts dagster.yaml there. Chart
concurrency / runRetries / runMonitoring mirror the same numbers.
Contributor gotcha (from __future__ import annotations)
Modules that define a class inheriting Config or
ConfigurableResource must not contain
from __future__ import annotations. Dagster 1.13.13 then raises
DagsterInvalidConfigDefinitionError with
'WorkflowConfig' cannot be resolved. Regression:
tests/dagster/test_config_annotations.py.
definitions.py itself may use postponed annotations because it does
not declare those subclasses.
Pools
Vendor-call assets declare pool="vendor:<service>" (for example
vendor:alphavantage, vendor:airbyte, vendor:datahub). Intended
named limits 1 / 2 are comments + VENDOR_POOL_LIMITS only
(vendor:alphavantage: 1, vendor:airbyte: 2, vendor:datahub: 2 in
yaml comments and
tests/dagster/test_concurrency_pools.py).
1.13.13 cannot store named pool limits. Runtime enforcement is instance
default_limit 8. Do not treat the named 1/2 figures as live caps.
Retries and run-failure
run_retries in yaml retries a failed run up to three times.
alphaswarm_run_failure then emits a credential-safe progress frame
through _progress.emit ({task_id, stage, message, timestamp, **extras}). It never includes exception text —
failure_event.message can carry secrets.
Freshness
build_last_update_freshness_checks covers ingest assets selected by
INGEST_SCHEDULE_NAMES (daily full refresh, weekday regulatory, daily
time-partitioned sources, daily Alpha Vantage partition, hourly Alpha
Vantage delta). Entity keys are excluded.
The window is the tightest cron gap per asset. The weekday
regulatory cron (0 4 * * 1-5) has 24h weekday gaps and a 72h
Friday-to-Monday gap; the builder takes the tightest interval, so that
ingest is 24h, not 72h. Daily ingest crons are 24h. Alpha Vantage
assets that also sit on 20 * * * * get a 1-hour window. There are no
separate freshness windows for profiling or DataHub sync.
One automation condition
AutomationCondition.eager() is attached only to
alphavantage_intraday_datahub_update. Sensor
alphaswarm_automation_condition_sensor targets that asset
(default_status=RUNNING).
Sandbox
Per-session interactive isolation (tempdir, Redis prefix, ContextVar endpoint overrides) is documented on Dagster sandbox.
Non-goals
- No
dg/create-dagstermigration. - The
alphaswarm_orchestrationDagster adapter is untouched. partitions.pystays unused.- No Dagster Plus.
- No multi-code-location deployment that loads both graphs in one instance.
Incident triage for Airbyte + Dagster still starts at the DataOps Airbyte and Dagster runbook.