Skip to content

Concurrency and Rate Limiting

Ductwork Pro lets you bound how many jobs may run against a shared resource at once, and how fast they may start. Steps declare which resource they consume with a limit_by: key in the pipeline definition. Your ductwork.yml declares what that resource permits: a concurrency ceiling, one or more rate limits, or both. Enforcement happens at claim time, so a limited step is never picked up by a worker unless the resource has room for it.

Limits are global. They span every worker thread, every process, and every host in your fleet, and they span every step and pipeline that shares the key. A limit of five means the count never reaches six.


Most limits exist because something downstream breaks when you exceed them:

  • Vendor connection caps: an API that permits five concurrent connections per account
  • Vendor throughput caps: a published limit like “1 request per second, 2,000 per hour”
  • Shared infrastructure: a legacy database, an SFTP server, or a queue that falls over under a fan-out
  • Cost control: metered services where an unbounded expand turns into an unbounded bill

Ductwork already has job_worker.count, but that is a per-pipeline, per-process thread count. It does not compose across a fleet, and it cannot protect a resource that three different steps in three different pipelines all call. A limit_by: key is a single pool that does.


Add limit_by: to any transition in your pipeline definition. It accepts a string or symbol naming the resource:

class SyncOrdersPipeline < Ductwork::Pipeline
define do |pipeline|
pipeline.start(FetchPendingOrders)
.expand(to: PushOrderToVendor, limit_by: :vendor_api)
.collapse(into: RecordSyncResults)
end
end

limit_by: is available on every transition: start, chain, expand, divide, divert, combine, converge, collapse, and dampen. It combines freely with delay: and timeout::

pipeline.chain(to: ChargeCard, delay: 5.minutes, timeout: 30.seconds, limit_by: :stripe)

Steps that share a key share a pool. If PushOrderToVendor and RefreshVendorCatalog both declare limit_by: :vendor_api, they draw from the same concurrency slots and the same rate budget, even when they live in different pipelines.

⚠️ Validation: The key must be a non-blank string or symbol. Anything else raises an ArgumentError at definition time.


Every key referenced by limit_by: needs an entry under job_worker.limits in ductwork.yml. An entry may declare concurrency, rate, or both:

default: &default
job_worker:
limits:
vendor_api:
concurrency: 5
rate:
- limit: 1
per: 1s
- limit: 2000
per: 1h
burst: 50
report_export:
concurrency: 2
bulk_import:
rate:
limit: 100
per: 1m
production:
<<: *default
staging:
<<: *default
job_worker:
limits:
vendor_api:
concurrency: 1

Because limits are ordinary configuration, they are scoped by environment like everything else in ductwork.yml. Staging can run concurrency: 1 against the same pipeline code that production runs at 5.

The maximum number of jobs stamped with this key that may be running at once, fleet-wide. A worker will not claim a job for this key while that many are already in flight.

One or more entries, each describing a sustained rate. rate accepts a single entry or a list:

OptionRequiredDescription
limityesHow many jobs may start per period
peryesThe period, as a whole number of seconds or a duration string
burstnoHow many jobs an idle key may start at once. Defaults to one second’s worth of limit, minimum 1

per accepts a bare integer (seconds) or a duration built from s, m, h, d, and w. Compound durations like 1h30m are fine.

per: 60 # sixty seconds
per: 1m # same thing
per: 1h30m # ninety minutes

Months and years are rejected because their length is approximate. Zero and sub-second periods are rejected too, but nothing is lost: per is a denominator, so any rate is expressible with a whole-second period. One job every N milliseconds is limit: 1000, per: Ns.

IntentConfig
100 per secondlimit: 100, per: 1s
One every 250mslimit: 4, per: 1s
One every 7mslimit: 1000, per: 7s
2,000 per hourlimit: 2000, per: 1h

Vendors routinely publish more than one limit, such as “1 per second and 2,000 per hour”. Neither entry can stand in for the other, so list both. A job starts only when every entry has budget for it.

Two rules are applied at boot to keep the list honest:

  • Dominated entries are dropped. If one entry is at least as strict as another on both rate and burst, the looser one is dead config and is silently ignored. 1 per 1s alongside 5000 per 1h keeps only the per-second entry, since 1/s is 3,600/h sustained.
  • Duplicate periods raise. Two entries under one key with the same per are rejected. That pair is always a slower bucket with a smaller burst written twice.

Limits use two new tables and two new columns. If you installed Ductwork Pro before this feature landed, run the Pro install generator again and migrate:

Terminal window
❯ bin/rails generate ductwork:pro:install
❯ bin/rails db:migrate

The generator only creates migrations you don’t already have. If you configure a limit and the tables are missing, Ductwork raises at boot naming them rather than failing on the first claim.


bin/ductwork start validates limits before any worker starts, so a typo fails the boot instead of quietly uncorking a vendor API:

ConditionResult
A limit_by: key has no entry under job_worker.limitsRaises, naming the key and the pipelines that reference it
An entry declares neither concurrency nor rateRaises. An empty entry is never silently unlimited
Two rate entries share a periodRaises
A period uses months or years, is zero, or can’t be parsedRaises
Tables are missing but limits are configuredRaises, naming the tables
An entry exists with no limit_by: reference anywhereWarns only (see Retiring a key)

The validator only sees pipelines that have been loaded. It runs after Ductwork’s own definition validation, which eager loads them, so this holds under the normal CLI boot.


Every worker poll asks one question: of the keys that have pending work and are permitted to run right now, which holds the oldest job? The worker claims from that one. Unkeyed work competes as one more lane in that comparison, so when nothing is saturated, the system still processes jobs in global first-in, first-out order.

For a keyed claim, the check happens under a per-key lock, in the same transaction as the claim itself:

  1. Lock the key’s header row, so concurrent workers on the same key serialize
  2. Count in-flight jobs for the key. At or over concurrency? Release and skip
  3. Refill every rate bucket by elapsed time and test for a token. Any bucket short? Release and skip
  4. Claim the oldest pending job for the key, then spend a token from every bucket, then commit

Because the count and the token test happen under the lock, two workers can’t both read “four in flight” against a limit of five and both claim. Enforcement is exact on every supported database, including SQLite.

The in-flight number is a live count of executions that are claimed and not yet completed. There is no stored counter to release. When a job finishes, errors, crashes, times out, or is swept by the reaper as an orphan, its slot comes back on its own. A leaked slot is impossible by construction.

A concurrency-saturated key is simply retried on the next poll, since there’s no way to know when a running job will finish.

Each rate entry is a token bucket: limit / per tokens refill continuously per second, up to a capacity of burst. A job spends one token from every bucket on its key. Tokens are stored in the database, so restarts don’t reset them, and every worker in the fleet draws from the same budget.

A few consequences worth knowing:

  • New buckets start empty. The first job on a freshly configured key waits roughly per / limit seconds. At the default burst that’s sub-second; at 2000 per 1h it’s about 1.8 seconds. Seeding full would hand a full burst to a key with no history.
  • Burst defaults small on purpose. Most token bucket implementations default capacity to limit, which would let an idle 2000 per 1h key fire 2,000 requests in one instant. Ductwork defaults to one second’s worth. Raise burst: explicitly if your vendor tolerates it. The worst case over any trailing window is burst + limit.
  • Config changes apply on restart, without a migration. limit and burst are read from config on every refill and never stored, so lowering burst from 50 to 5 takes effect immediately: banked tokens above the new capacity are clamped away on the next poll.
  • Rate-blocked keys don’t spin. When a spend or a refusal leaves a bucket short, the key is stamped with the exact time it will have a token again. Every worker skips the key until then, and the poll goes to other keys or unkeyed work instead.

The limit governs execution claims, not outbound requests. That is deliberate: making it exact at the request site would require a mid-job API you’d have to call yourself. The consequences:

  • Every retry spends a token. A step failing in a tight retry loop consumes rate budget at the retry rate. Each attempt is a real outbound request, so this is correct, but it’s the opposite of what most people assume.
  • A crash after claim spends a token. If a worker dies between claiming and calling out, the reaper re-queues the job and the retry spends another token.
  • An early return spends a token. A guard clause that skips the external call still consumed a claim.

Concurrency slots are always released correctly in all of these cases; only rate budget is affected.


Adding limit_by: does not limit runs already in flight. Like timeout: and delay:, the key is read from the definition snapshot taken when the run was triggered. Runs triggered after the deploy are limited from their first step. Runs already in progress keep creating unlimited work until they reach a terminal state, which for a dampened pipeline can be a long time. Every stamped row is enforced exactly; this is a coverage gap during rollout, not an overshoot.

Changing a limit’s value applies to everything on restart. Only the reference is versioned with the run. The values are resolved against ductwork.yml at claim time, so raising concurrency: 5 to 10 takes effect for every run, in-flight or new, as soon as workers restart.

Mid-deploy skew is self-resolving. While a rolling deploy has workers on two versions of ductwork.yml, they enforce different values against the same key. Every claim is still exact under the lock, nothing exceeds the looser of the two values, and it resolves when the rollout finishes.

Remove a key in two deploys:

  1. Remove the limit_by: reference from the DSL and deploy. Boot warns about the unreferenced config entry but does not fail. New runs stop creating keyed work; in-flight runs drain.
  2. Once keyed work has drained, remove the entry from job_worker.limits and deploy.

If you remove both at once, work already stamped with the key is not stranded, but it runs unlimited with a warning logged per claim ("Claiming a limit key with no configuration entry unlimited"). Boot-time validation fails closed; the claim path fails open, because a row that already exists can only be run or left forever.


Limits + Delays: A step can be both delayed and limited. The delay governs when the job becomes eligible; the limit governs whether a worker may claim it once it is. Delayed work never blocks a key while it waits.

Limits + Timeouts: A timed-out job’s execution is completed, so its concurrency slot is released immediately. Its retry, if any, spends another rate token.

Limits + Interruptible Advancement: A large expand onto a limited step works exactly as you’d hope. Millions of expanded steps are materialized in batches, and the limit meters how fast workers pick them up.