Saltar al contenido principal

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).

LocationModuleHow it is launchedWhat it owns
Monolithalphaswarm.dagster.definitionsCompose 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
Platformpipelines.dagster_user_code.definitionsHelm 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.

KindCountNames
Jobs12full_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
Schedules10daily_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 * * * *
Sensors3pipeline_manifests_changed, alphaswarm_run_failure, alphaswarm_automation_condition_sensor
Asset checks6datahub_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 checks2last-update checks grouped by tightest ingest-cron interval
Assets36alphaswarm_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.

FieldValue
run_retries.max_retries3
run_monitoring.free_slots_after_run_end_seconds300
concurrency.pools.default_limit8
concurrency.runs.max_concurrent_runs8

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-dagster migration.
  • The alphaswarm_orchestration Dagster adapter is untouched.
  • partitions.py stays 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.