Skip to content

Dagster Orchestration

Thomas Maerz edited this page Oct 4, 2026 · 1 revision

Dagster Orchestration

Workspace assets

Each managed workspace receives assets representing:

  • raw Slackdump extraction;
  • attachment reconciliation;
  • canonical DuckDB ingestion.

Blocking checks include:

  • archive_valid
  • attachment_integrity
  • canonical_integrity

Jobs

Slackpipe generates per-workspace initial and incremental jobs and coordinated all-workspace extract/canonical jobs. Initial jobs enforce the full archive contract; incremental jobs can use bounded stale rescans.

Schedules

Generated schedules default to stopped. This prevents newly discovered workspaces from beginning extraction before an operator reviews credentials, storage, and initial-run expectations.

Sensors

  • slackpipe_new_workspace_sensor detects newly configured workspaces.
  • slackpipe_rollout_coordinator advances serialized extraction, attachment, and canonical stages in configured order.
  • slackpipe_run_failure_metrics publishes bounded failure observations.

Concurrency

Dagster uses a queued run coordinator, per-operation pools, phase-specific tag limits, and a serial tag for workflows that must not overlap. Advisory locks remain the final process-level guard around raw and canonical mutation.

Shared deployment with Slackquery

The included workspace configuration can load both slackpipe and slackquery gRPC code locations. They share Dagster infrastructure but not ownership: Slackpipe writes canonical data; Slackquery writes only its own state/artifacts.

Clone this wiki locally