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.
Why limit concurrency and rate?
Section titled “Why limit concurrency and rate?”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
expandturns 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.
Declaring a key in the DSL
Section titled “Declaring a key in the DSL”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) endendlimit_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
ArgumentErrorat definition time.
Configuring limits
Section titled “Configuring limits”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: 1Because 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.
concurrency
Section titled “concurrency”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:
| Option | Required | Description |
|---|---|---|
limit | yes | How many jobs may start per period |
per | yes | The period, as a whole number of seconds or a duration string |
burst | no | How 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 secondsper: 1m # same thingper: 1h30m # ninety minutesMonths 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.
| Intent | Config |
|---|---|
| 100 per second | limit: 100, per: 1s |
| One every 250ms | limit: 4, per: 1s |
| One every 7ms | limit: 1000, per: 7s |
| 2,000 per hour | limit: 2000, per: 1h |
Multiple rate entries
Section titled “Multiple rate entries”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 1salongside5000 per 1hkeeps only the per-second entry, since 1/s is 3,600/h sustained. - Duplicate periods raise. Two entries under one key with the same
perare rejected. That pair is always a slower bucket with a smaller burst written twice.
Installation
Section titled “Installation”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:
❯ bin/rails generate ductwork:pro:install❯ bin/rails db:migrateThe 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.
Boot-time validation
Section titled “Boot-time validation”bin/ductwork start validates limits before any worker starts, so a typo fails the boot instead of quietly uncorking a vendor API:
| Condition | Result |
|---|---|
A limit_by: key has no entry under job_worker.limits | Raises, naming the key and the pipelines that reference it |
An entry declares neither concurrency nor rate | Raises. An empty entry is never silently unlimited |
| Two rate entries share a period | Raises |
| A period uses months or years, is zero, or can’t be parsed | Raises |
| Tables are missing but limits are configured | Raises, naming the tables |
An entry exists with no limit_by: reference anywhere | Warns 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.
How enforcement works
Section titled “How enforcement works”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:
- Lock the key’s header row, so concurrent workers on the same key serialize
- Count in-flight jobs for the key. At or over
concurrency? Release and skip - Refill every rate bucket by elapsed time and test for a token. Any bucket short? Release and skip
- 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.
Concurrency is derived, never counted
Section titled “Concurrency is derived, never counted”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.
Rate is a token bucket
Section titled “Rate is a token bucket”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 / limitseconds. At the default burst that’s sub-second; at2000 per 1hit’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 idle2000 per 1hkey fire 2,000 requests in one instant. Ductwork defaults to one second’s worth. Raiseburst:explicitly if your vendor tolerates it. The worst case over any trailing window isburst + limit. - Config changes apply on restart, without a migration.
limitandburstare read from config on every refill and never stored, so loweringburstfrom 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.
What the limit governs
Section titled “What the limit governs”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.
Rollout and change management
Section titled “Rollout and change management”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.
Retiring a key
Section titled “Retiring a key”Remove a key in two deploys:
- 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. - Once keyed work has drained, remove the entry from
job_worker.limitsand 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.
Combining with other features
Section titled “Combining with other features”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.