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.
Ancestors
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.
Ancestors
_DAGChildOkResult
class _DAGChildOkResult():The flow run finished. Its outputs are already in the cache.
Ancestors
_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 topod_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-levelRunrow 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 sametask_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