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 openPrefectClient.deployment_id: The deployment whose queue to advance.exclude_flow_run_id: Seeruns_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 openSyncPrefectClient.deployment_id: The deployment whose queue to advance.exclude_flow_run_id: Seeruns_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;Noneis 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 stillRunningas it does.
Returns The parked runs to promote, in order.