skip to content

In Prefect, how would you cap concurrent calls that many flows make to one shared API?

level: principalimportance: should knowfreq 30%

answer

  1. the limit belongs to the resource, not the flow
  2. one counter every worker consults
  3. hold the slot for the call, not the task
  4. a killed run may not give the slot back

basics

~20 s

Put the cap in Prefect's concurrency system rather than in each flow: a task tag concurrency limit throttles every task run carrying that tag across the workspace, and a global concurrency limit with the concurrency or rate_limit context manager guards arbitrary code blocks.

solid answer

~50 s

Per-flow throttling does not compose — five flows each allowing four in-flight calls is twenty. Prefect gives you two shared mechanisms. **Tag-based task concurrency limits** attach a maximum to a tag: tag every task that touches the API with `vendor-api`, set the limit once, and Prefect will only let that many tagged task runs be Running at a time across the whole workspace. **Global concurrency limits** are named counters you take explicitly with the `concurrency("vendor-api", occupy=1)` context manager around just the call, plus `rate_limit(...)` when the constraint is requests per second rather than simultaneous requests. Choose the tag when the whole task is the unit of work and the global limit when only part of a task should hold the slot. Then plan the operational edges: a killed run can leave a slot occupied until it is reset, holding a slot while awaiting other work that needs the same limit deadlocks, and throttling never removes the need to retry a 429 with backoff.

code

python · 13 lines
python
from prefect import flow, task
from prefect.concurrency.sync import concurrency, rate_limit

@task(retries=3, retry_delay_seconds=10, timeout_seconds=120)
def fetch_customer(customer_id: str) -> dict:
    payload = prepare(customer_id)          # no slot held
    rate_limit("vendor-api-rps")            # pace requests per second
    with concurrency("vendor-api", occupy=1):
        return call_vendor(payload)         # slot held only here

@flow
def sync_all(ids: list[str]):
    fetch_customer.map(ids)                 # fan-out stays bounded

go deeper

for a junior

Know that Prefect can limit how many tagged task runs execute at once, and that the limit is configured centrally rather than written into each flow's code.

for a middle

Explain the two mechanisms — tag-based task concurrency limits versus named global limits taken with the concurrency context manager — and that slots are released when a run leaves the Running state.

for a senior

Show the operating detail: slots leaked by killed runs, timeouts so starved runs fail visibly instead of hanging, and pairing a throttle with retries and backoff because a shared cap never removes rate-limit errors.

for a principal

Own the policy — which layer each constraint belongs to, weighting with occupy, avoiding self-deadlock, and the governance of a workspace-wide number that silently arbitrates throughput between teams.

## Why per-flow limits are the wrong layer The constraint you are protecting belongs to the **API**, not to any one flow. If each flow limits itself to four concurrent calls, the vendor sees four times however many flows happen to run at once — a number that changes every time someone adds a deployment. Any cap that is going to hold has to be expressed once, somewhere both flows can see. In Prefect that is the concurrency system, which stores counters server-side so every worker, pod and flow run consults the same accounting. ## Mechanism one: tag-based task concurrency limits Prefect task runs carry tags, and a tag can be given a maximum number of simultaneously Running task runs: ```python @task(tags=["vendor-api"], retries=3, retry_delay_seconds=10) def fetch(customer_id: str): ... ``` With a limit of 5 on `vendor-api`, at most five task runs bearing that tag execute at once, no matter which flow submitted them or which worker picked them up. Runs beyond the limit wait for a slot rather than executing. The slot is held for as long as the task run is in a Running state and released when it leaves. This is the right tool when the **task is the unit of contention** — one task, one call, one slot. It is a blunt instrument when a task does an API call and then thirty seconds of local work: the slot is held for the whole task, so your effective concurrency against the API is lower than the number you set. ## Mechanism two: global concurrency limits and rate limits Global concurrency limits are named, server-side counters you acquire explicitly around any block of code: ```python from prefect.concurrency.sync import concurrency, rate_limit @task def fetch(customer_id: str): rows = prepare(customer_id) # unthrottled with concurrency("vendor-api", occupy=1): payload = call_vendor(rows) # only this holds a slot return transform(payload) # unthrottled ``` Two things this buys you over tags. First, precision: only the contended section occupies the slot. Second, weighting — `occupy` lets an expensive operation take several slots, so a bulk call and a single-row call can share one budget honestly. `rate_limit(...)` addresses the other shape of vendor constraint: not "how many at once" but "how many per second", pacing releases rather than counting concurrent holders. An async variant of the same API exists for async flows. ## Mechanism three: the infrastructure layer Work pools and deployments can also cap how many runs execute concurrently. That is a coarser, infrastructure-shaped control — it protects your cluster and your budget rather than the vendor — and it is the wrong knob for this problem, because it throttles whole flow runs regardless of whether they touch the API at all. ## Choosing, as a design decision - Contention is the entire task, and the tag is a natural label for it → **tag limit**. - Contention is one call inside a larger task, or different calls cost different amounts → **global limit with `occupy`**. - The vendor publishes requests-per-second rather than concurrency → **rate limit**. - You are protecting your own compute or spend → **work pool / deployment concurrency**. Often the honest answer is two of them: a rate limit for the vendor's published quota and a work-pool cap so a backfill cannot flood your cluster. ## The operational edges a lead is expected to name **Leaked slots.** A run killed by its infrastructure never leaves Running cleanly, so its slot can stay occupied. The queue then drains more slowly than the limit implies, or stalls entirely. Know that slots can be reset, and monitor for a limit sitting pinned at its maximum with nothing actually running. **Self-deadlock.** If a flow acquires a slot and then waits on subflows or tasks that need the same limit, the parent holds the resource its children need and nothing progresses. Acquire around leaves, not around orchestration. **Queueing is not free.** Waiting task runs still occupy a worker's attention or sit in a waiting state; a very low limit with a very wide fan-out turns into a long tail of slow runs and timeouts that look like failures. Set `timeout_seconds` deliberately so a starved run fails visibly instead of hanging. **Throttling is not error handling.** A shared cap reduces the chance of a 429 but does not eliminate it — other consumers exist, and vendor limits change. Keep `retries` with backoff on the calling task, and treat the throttle as capacity shaping rather than correctness. **Governance.** Because these limits are workspace-level, they are shared state: one team raising a limit affects another team's throughput. Decide who owns each named limit, document what constraint it encodes and where the number came from, and review it when the vendor contract changes — an undocumented number that nobody dares touch is the usual end state.

  • When is a tag limit the wrong choice compared with a global concurrency limit?
    When only part of the task contends. A tag holds the slot for the whole task run, so a task that calls the API for two seconds and then transforms for thirty wastes most of its slot. A global limit acquired with a context manager around just the call keeps the throttle tight, and occupy lets heavier calls take more than one slot.
  • How can a concurrency limit stall even when nothing is running?
    Slots are released when a run leaves the Running state, so a run killed by its infrastructure can leave its slot occupied. The limit then sits pinned at its maximum with no real work behind it and the queue never drains. Watch for that shape, know how to reset the limit's active slots, and prefer timeouts so hung runs terminate.
  • Can holding a slot deadlock a flow against itself?
    Yes. If a flow acquires a slot and then waits on tasks or subflows that need the same limit, the parent is holding exactly what its children need and nothing progresses. Acquire around the leaf operation that actually contends, never around an orchestration step that waits on other runs sharing the counter.

saying these in an interview costs you the question

  • Sets the concurrency cap separately inside each flow
  • Assumes throttling removes the need to retry 429s
  • Holds a slot across an entire long task
  • Wraps orchestration that waits on same-limited work
  • Treats a workspace-wide limit as one team's private setting

context