Skip to main content

executor

Executor: runs a RunnableDAG inside a Prefect flow.

This is the runtime heart of the DAG system. It takes a fully parsed and validated RunnableDAG (a BackgroundDAG or InteractiveDAG) and executes every step in declaration order, wiring inputs from prior step results or context providers.

Step execution​

Every step runs off the event loop on a daemon thread, which keeps the (synchronous) task function clear of the loop's default thread-pool executor - whose non-daemon workers would hold the process open at shutdown. Results stay in-process — no Prefect serialisation boundary is crossed — so private attributes such as _accessor on cache-backed result types are preserved for downstream steps.

Nodes marked parallel=True are additionally run concurrently: they are tracked as asyncio.Task objects in a pending dict and awaited lazily when a downstream step declares a FromRef dependency, or flushed at the end of the loop. Sequential nodes are awaited immediately, so declaration order holds.

Sequential steps used to be called directly on the event loop. Because the pod runs its background DAGs in-process, on the same loop that supervises the mailbox listener and dispatches its handlers, a long step — an hour has been observed — left the pod unable to see listener failures or run queued message handlers until it finished.

A background run's context assembly is off the loop for the same reason, and on the same kind of thread: build_background_flow awaits open_run rather than calling it, because materialising the datasource and planning the waves is work of that same order.

Runtime parameters​

Steps often require runtime objects (datasource, cache, task_hash, …) that are not declared in the YAML. Pass them via DAGRunContext.runtime_params; they are merged into every step's kwargs before YAML-resolved inputs, so YAML refs can override them when needed.

Public API​

execute_dag(dag, run_ctx) Runs the DAG. Caller must already be inside a Prefect @flow context; raises RuntimeError otherwise.

run_dag(dag, run_ctx) Ensures a Prefect flow context exists, then delegates to execute_dag. Safe to call from anywhere.

build_flow(dag, run_ctx) Returns a named Prefect @flow callable for deployment scenarios.

execute_background_dag(dag, *, task_hash, project_id, open_run, ...) Runs one background DAG from inside whichever flow run already owns the work — the served background_dag_runtime deployment, which must not nest a second tagged flow run inside its own.

build_background_flow(dag, *, task_hash, project_id, open_run, concurrency_slot=None, slot_timeout_seconds=None, on_slot_acquired=None) Returns a @flow callable that opens the run inside the flow body, so a background run's flow run is the outermost boundary of its work. A rerun trigger also passes a concurrency_slot, acquired around the body, and an on_slot_acquired hook that fires only once the slot is held.

Module​

Functions​

build_background_flow​

def build_background_flow(    dag: BackgroundDAG,    *,    task_hash: str,    project_id: str,    open_run: Callable[[], DAGRunContext | None],    concurrency_slot: str | None = None,    slot_timeout_seconds: float | None = None,    on_slot_acquired: Callable[[], None] | None = None,    on_completed: Callable[[], None] | None = None,    on_success: Callable[[], None] | None = None,) ‑> collections.abc.Callable[[], collections.abc.Coroutine[typing.Any, typing.Any, dict[str, typing.Any] | prefect.client.schemas.objects.State]]:

Return a @flow callable that opens its own run, inside the flow.

The background counterpart to build_flow, and the shape that makes "a Prefect flow run exists" a total invariant for background work. build_flow needs a DAGRunContext in hand, which means the run row, the datasource materialisation and the model resolution all happen before anything is observable; a crash in that window leaves a running cache row and no flow run to explain it. Here open_run is deferred into the flow body, so all of it is inside the boundary: a setup failure is a failed flow run rather than a log line and a dropped trigger.

The duplicate guard lives here too, for the same reason — it is now Prefect's flow-run state, not the cache's running row, that decides whether a second run may start (see bitfount.runtimes.dedup.is_background_flow_run_active, which is stale-aware so a dead run cannot wedge the lineage). A duplicate ends as an immediately-terminal SKIPPED_STATE_NAME run naming the live run it duplicates, which is visible in the Prefect UI where the old early return None was visible nowhere.

Arguments

  • dag: The validated background DAG to execute.
  • task_hash: The (pod, datasource) task hash, for the duplicate guard.
  • project_id: Project the run belongs to, for the duplicate guard.
  • open_run: Deferred run-context assembly — typically functools.partial(build_background_run_context, ...). Called once, inside the flow body, after the guard has passed. None means the caller declined the run for its own reasons, and is reported the same way a duplicate is.
  • concurrency_slot: Name of a Prefect global concurrency limit to hold for this run's execution, or None to run unbounded. Rerun triggers pass one; a dataset-project link and a recovery replay do not, so work a user is waiting on never queues behind a discretionary rerun that may hold its slot for hours (Prefect global limits have no priority and no preemption).
  • slot_timeout_seconds: How long to wait for the slot before ending the run skipped. None waits indefinitely.
  • on_slot_acquired: Called once, from inside the flow body, after the slot is held and before open_run. The hook exists so a caller can record "this lineage is now running" at the only moment that is
  • true: a rerun that loses the race for the slot ends skipped without it, and stays due for the next tick. Called in the same position on the unslotted path, so a link or a recovery replay sees identical ordering.

Raising from it fails the flow run and nothing executes. That is the cheaper of two bad answers rather than a good one: the lineage is left un-stamped, so its scheduler finds it due again on the next pass and this repeats until the write succeeds — one failed run per tick, which is loud and costs nothing, where executing the DAG anyway would redo the whole lineage per tick instead. A hook that keeps failing is a broken cache, and the log line to look for.

  • on_completed: Called once, from inside the flow body, after the DAG has executed without raising and without abandoning a wave. What it records is "this lineage carried what it started", which is true of a budgeted refresh that deliberately carried its allowance and stopped, and false of a run that gave up on a wave.

The line between the two hooks is deferring work versus failing at it. Deferring is a scheduling decision and must not look like a fault, or a healthy multi-run sweep would resemble a failing lineage within a few nights — which is why this is not gated on covering the datasource the way on_success is. Abandoning a wave is a fault, and must not be cleared here: the streak a caller resets from this hook is the only thing that stands a lineage down when a wave fails permanently, and a refresh continuation is driven by the ledger still being stale. Clear it on a run that abandoned a wave and that wave is retried every poll tick until the grace window closes, holding the pod-wide rerun slot each time.

Not called on either Skipped exit, nor on a raising one.

  • on_success: Called once, from inside the flow body, after the DAG has executed without raising. Its counterpart to on_slot_acquired: that one records that a lineage started, this one records what a lineage covered, and the two must not be collapsed. A caller banking "the inventory was this big when I finished" from the start of a run would bank it for files the run never reached if it then died, and those files would stop looking outstanding.

Not called on either Skipped exit — a duplicate and a declined run both covered nothing — nor on a raising one, which propagates and fails the flow run as before.

Returns An async Prefect @flow callable with no required arguments. It returns the step results, or a terminal State when nothing ran.

build_flow​

def build_flow(    dag: RunnableDAG,    run_ctx: DAGRunContext,) ‑> collections.abc.Callable[[], collections.abc.Coroutine[typing.Any, typing.Any, dict[str, typing.Any]]]:

Return a named Prefect @flow callable for dag.

Useful for deployment scenarios where a first-class flow object is needed. For most runtime cases prefer run_dag.

Note that the flow's name is dag.name, which is shared by every project and datasource using the same template — it does not identify a run. Flow runs that need to be findable again are tagged at call time by _run_dag_flow_and_flush; see bitfount.runtimes.dedup.background_dag_tags.

Arguments

  • dag: A validated RunnableDAG (BackgroundDAG or InteractiveDAG).
  • run_ctx: Runtime context passed through to execute_dag.

Returns An async Prefect @flow callable with no required arguments.

execute_background_dag​

async def execute_background_dag(    dag: BackgroundDAG,    *,    task_hash: str,    project_id: str,    open_run: Callable[[], DAGRunContext | None],    concurrency_slot: str | None = None,    slot_timeout_seconds: float | None = None,    on_slot_acquired: Callable[[], None] | None = None,    on_completed: Callable[[], None] | None = None,    on_success: Callable[[], None] | None = None,) ‑> dict[str, typing.Any] | prefect.client.schemas.objects.State:

Run one background DAG, from inside whichever flow run owns the work.

The body build_background_flow wraps, factored out so that a caller which already is the flow run can use it without nesting a second one. That caller is the served background_dag_runtime deployment: its own flow run carries this lineage's tags, so a subflow carrying them too would be seen by the duplicate guard as a live run of the same lineage and skip the work it was launched to do.

Guards, takes the slot, opens the run, then executes the DAG. In that order, and the order is the contract: nothing expensive is assembled until a slot is held, and on_slot_acquired fires only once it is — so a run that never got one recorded nothing either.

Every argument means what it means on build_background_flow; see there.

Arguments

  • dag: The validated background DAG to execute.
  • task_hash: The (pod, datasource) task hash, for the duplicate guard.
  • project_id: Project the run belongs to, for the duplicate guard.
  • open_run: Deferred run-context assembly, called inside the guard.
  • concurrency_slot: Prefect global concurrency limit to hold, if any.
  • slot_timeout_seconds: How long to wait for that slot.
  • on_slot_acquired: Called once the slot is held, before open_run.
  • on_completed: Called after the DAG executed without raising and without abandoning a wave.
  • on_success: Called after the DAG executed without raising and covered the datasource.

Returns The step results, or a terminal State when nothing ran.

execute_dag​

async def execute_dag(dag: RunnableDAG, run_ctx: DAGRunContext) ‑> dict[str, typing.Any]:

Run dag inside an already-active Prefect flow context.

Raises RuntimeError if no flow context is active — use run_dag when you cannot guarantee one.

Arguments

  • dag: A validated RunnableDAG (BackgroundDAG or InteractiveDAG).
  • run_ctx: Runtime context: providers, runtime params, reporter.

Returns Mapping of step name → step result for every executed step.

Raises

  • RuntimeError: If called outside a Prefect flow context.
  • Exception: Re-raises any step exception after calling run_ctx.lifecycle_notifier.on_failure.

run_dag​

async def run_dag(dag: RunnableDAG, run_ctx: DAGRunContext) ‑> dict[str, typing.Any]:

Run dag, creating a Prefect flow context if one is not already active.

Safe to call from anywhere — no existing flow context required.

Arguments

  • dag: A validated RunnableDAG (BackgroundDAG or InteractiveDAG).
  • run_ctx: Runtime context: providers, runtime params, reporter.

Returns Mapping of step name → step result for every executed step.

run_flow_resiliently​

async def run_flow_resiliently(flow_fn: Callable[[], Awaitable[Any]]) ‑> Any:

Run a Prefect flow, retrying the transient Windows start-time 409.

Guards against the Windows-only flow_run_state timestamp collision (see the module note above): a bounded retry re-creates the flow run a clock-tick later so its Pending/Running state timestamps differ. Only a 409 is retried; any other error — and the final 409 after exhausting retries — is re-raised unchanged. Safe against re-running work: the collision happens at begin_run before any step executes (a real DAG does far more than one clock tick of work, so the Running->Completed transition never shares a tick).

Each retry creates a new flow run, so the run it gave up on is left behind at PENDING, carrying this run's lineage tags. Two consequences, handled separately because they need different remedies:

  • The duplicate guard must not treat a PENDING run as live, or the retry would find the run it just abandoned and skip itself as a duplicate — see bitfount.runtimes.dedup.BACKGROUND_SUBMITTED_STATE_TYPES, which counts SCHEDULED and RUNNING and excludes PENDING for exactly this reason. That is correctness, and it does not depend on the cleanup below succeeding.
  • The abandoned run is ended as CANCELLED here (_abandon_flow_run), so it does not sit PENDING forever misrepresenting the pod's history. That is tidiness, and it is best-effort.