Skip to main content

dag_process

Run a background v9 DAG in a dedicated child process.

The pod must stay responsive while a DAG runs: its heartbeat, message handlers and liveness tickers share one event loop, and a step that blocks that loop takes the pod's datasets offline on the Hub. A separate process makes that impossible by construction.

A background run needs no IPC beyond a terminal status. Its results are written to the cache, which is the only sink, and its LifecycleNotifier carries no mailbox, so no lifecycle traffic crosses the boundary.

State that cannot cross a process boundary is rebuilt here: the Hub session, the cache handle, the cache's encryption binding and the pod's hooks. The datasource crosses as a pickled DatasourceContainer.

The child works from a datasource snapshot taken when it was spawned, so a config reload mid-run does not reach it. A run reads a consistent view of its datasource.

Module​

Functions​

_run_background_dag_in_child​

def _run_background_dag_in_child(    config: _DAGProcessConfig,    result_queue: multiprocessing.Queue[Any],    heartbeat: Heartbeat | None,) ‑> None:

Child entry point: run one background DAG and report how it ended.

Must stay importable at module level — spawn pickles it by reference.

Arguments

  • config: This run's configuration.
  • result_queue: Where the terminal status goes.
  • heartbeat: Shared timestamp the parent ages to tell a working child from a wedged one.

Classes​

_DAGChildCacheUnusableResult​

class _DAGChildCacheUnusableResult(error_message: str):

The background cache could not be opened or migrated.

Reported rather than acted on: discarding the cache is the parent's decision, and the parent owns the flag that stands background work down.

Arguments

  • error_message: What made the cache unusable.

Variables​

  • static error_message : str

_DAGChildErrorResult​

class _DAGChildErrorResult(    error_message: str, exception_type: str, traceback_summary: str,):

The run failed.

Carries strings rather than the exception: an exception that does not pickle would otherwise replace the real fault with a pickling error on the way back.

Arguments

  • error_message: str() of the exception.
  • exception_type: Its class name.
  • traceback_summary: The formatted traceback.

Variables​

  • static error_message : str
  • static exception_type : str
  • static traceback_summary : str

_DAGChildOkResult​

class _DAGChildOkResult():

The flow run finished. Its outputs are already in the cache.

_DAGChildResultMessage​

class _DAGChildResultMessage():

Base for what the child puts on the result queue.

_DAGProcessConfig​

class _DAGProcessConfig(    pod_name: str,    pod_identifier: str,    pod_key_path: Path | None,    username: str | None,    hub_url: str,    secrets: APIKeys | RefreshableJWT | dict[SecretsUse, APIKeys | RefreshableJWT] | None,    flow_spec: FlowSpec,    ds_container: DatasourceContainer,    task_hash: str,    run_id: str,    project_id: str,    datasource_name: str,    monitor_task_id: str,    trigger: BackgroundTrigger,    is_rerun: bool,    run_tags: list[str],    concurrency_slot: str | None,    slot_timeout_seconds: float | None,    ehr_config: EHRConfig | None = None,    ehr_secrets: RefreshableJWT | None = None,    hook_factories: list[Callable[[], None]] = [],):

Everything the child needs to run one background DAG.

All fields are picklable: this crosses a spawn boundary. Live objects — the Hub session, the cache handle, the notifier — are rebuilt in the child from the primitives here rather than shipped.

Arguments

  • pod_name: Pod whose background cache this run writes to.
  • pod_identifier: username/pod_name, for Hub reporting and telemetry.
  • pod_key_path: Path to pod_rsa.pem, bound into the cache-encryption layer so encrypted columns can derive their key.
  • username: Hub username, used to rebuild the session and tag telemetry.
  • hub_url: Hub base URL.
  • secrets: Credentials the child authenticates with.
  • flow_spec: The parsed spec. The DAG is rebuilt from it in the child rather than pickled, so the child and the pod agree by construction.
  • ds_container: The datasource, its schema and its config, as a snapshot.
  • task_hash: The (pod, datasource) hash keying this run.
  • run_id: The id this run's DAG-level Run row is opened under. Minted by the parent, not the child: a child that dies mid-run cannot close its own rows, and the id is what lets the parent close this run's rows without touching a sibling's on the same task_hash.
  • project_id: Project the datasource is linked to.
  • datasource_name: Datasource this run processes.
  • monitor_task_id: Hub monitor task id to report under.
  • trigger: Why this run was started.
  • is_rerun: Whether to apply rerun semantics — hold the concurrency slot and stamp the lineage.
  • run_tags: Tags to attach to the Prefect flow run.
  • concurrency_slot: Prefect concurrency limit to occupy, if any.
  • slot_timeout_seconds: How long to wait for that slot.
  • ehr_config: EHR configuration, if the pod has one.
  • ehr_secrets: EHR credentials, if the pod has them.
  • hook_factories: Picklable factories that re-register hooks in the child.

Variables​

  • static concurrency_slot : str | None
  • static datasource_name : str
  • static ds_container : DatasourceContainer
  • static ehr_config : EHRConfig | None
  • static ehr_secrets : RefreshableJWT | None
  • static flow_spec : FlowSpec
  • static hook_factories : list[Callable[[], None]]
  • static hub_url : str
  • static is_rerun : bool
  • static monitor_task_id : str
  • static pod_identifier : str
  • static pod_key_path : Path | None
  • static pod_name : str
  • static project_id : str
  • static run_id : str
  • static run_tags : list[str]
  • static secrets : APIKeys | RefreshableJWT | dict[SecretsUse, APIKeys | RefreshableJWT] | None
  • static slot_timeout_seconds : float | None
  • static task_hash : str
  • static trigger : BackgroundTrigger
  • static username : str | None