Skip to main content

queueing

Keep queued metadata runs in submission order behind a concurrency limit.

A deployment with a concurrency_limit enqueues on collision: a run that finds no free slot is returned to Scheduled(AwaitingConcurrencySlot) thirty seconds out. The Runner is shared by every metadata deployment, so it admits more runs of one deployment than that deployment allows, and the ones bounced back cycle round-robin: whichever is due when a slot frees takes it, whatever order the runs were submitted in.

So a submitter parks runs instead: every run is scheduled PARK_DELAY in the future, outside the Runner's prefetch window, one millisecond apart in submission order, and tagged PARKED_TAG. Then fill_free_slots moves the earliest parked runs to now, as many as the deployment has free slots. Each run does the same as it finishes, so a slot is handed to the next run in order rather than contended for.

A missed hand-over (a run that is cancelled or killed skips its own) is not lost: the next run to finish, the next submission, a cancellation or the startup recovery fills the slot, and the nightly refresh fills whatever is free every night. Failing all of those a parked run starts on its own at its parked time.

Module​

Functions​

fill_free_slots​

async def fill_free_slots(    client: PrefectClient, deployment_id: UUID, exclude_flow_run_id: UUID | None = None,) ‑> int:

Move as many parked runs to now as the deployment has free slots.

Never raises: a failure leaves the parked runs to start at their parked time.

Arguments

  • client: An open PrefectClient.
  • deployment_id: The deployment whose queue to advance.
  • exclude_flow_run_id: See runs_to_promote.

Returns How many runs were promoted.

fill_free_slots_sync​

def fill_free_slots_sync(    client: SyncPrefectClient,    deployment_id: UUID,    exclude_flow_run_id: UUID | None = None,) ‑> int:

fill_free_slots for a synchronous caller, such as a sync flow.

Arguments

  • client: An open SyncPrefectClient.
  • deployment_id: The deployment whose queue to advance.
  • exclude_flow_run_id: See runs_to_promote.

Returns How many runs were promoted.

hand_over_slot​

def hand_over_slot() ‑> int:

From inside a finishing deployment run, promote the next parked run.

A no-op outside a deployment run. Never raises.

Returns How many runs were promoted.

limit_of​

def limit_of(deployment: object) ‑> int | None:

A deployment's concurrency limit, or None if it has none.

Read from global_concurrency_limit: concurrency_limit is deprecated and always None in Prefect 3.

Arguments

  • deployment: A Prefect deployment response.

parked_time​

def parked_time(base: datetime, position: int) ‑> datetime.datetime:

The parked scheduled time of the run at position in a submission.

Arguments

  • base: When the submission started.
  • position: The run's index in submission order.

runs_to_promote​

def runs_to_promote(    runs: Iterable[FlowRun],    limit: int | None,    now: datetime,    exclude_flow_run_id: UUID | None = None,) ‑> list[prefect.client.schemas.objects.FlowRun]:

Choose which parked runs to move to now, earliest parked first.

Arguments

  • runs: A deployment's live runs.
  • limit: The deployment's concurrency limit; None is unlimited.
  • now: The current time.
  • exclude_flow_run_id: A run to leave out of the count of occupied slots: the run handing its slot over, which is still Running as it does.

Returns The parked runs to promote, in order.