Skip to content

Limit concurrency to N tasks at a time snippet

Run a thousand tasks but only N at a time — the bounded pool that protects downstream services and memory.

Run a thousand tasks but only N at a time — the bounded pool that protects downstream services and memory. The trap is the naive batch loop: await Promise.all(chunk) waits for the SLOWEST member of each batch before the next starts (a straggler stalls the lane), and one rejection can orphan the rest. Correct shapes: N workers pulling a shared queue, or a semaphore acquired before and released after — and always propagate the first error.

Runnable recipe · 12 languages
Concurrency & Parallelismconcurrencyasyncsemaphoreworker-poolparallelismrate-limiting

Every language

12 implementations, copy-ready. One at a time with syntax highlighting, or all inline.

JSJavaScript
// p-limit's core, hand-rolled: an active counter + a queue of resolvers.
function limitConcurrency(n) {
  let active = 0;
  const queue = [];

  function next() {
    if (active >= n || queue.length === 0) return;
    active++;
    const { fn, resolve, reject } = queue.shift();
    fn().then(resolve, reject).finally(() => {
      active--;
      next(); // a slot just freed — admit the oldest waiter
    });
  }

  return (fn) =>
    new Promise((resolve, reject) => {
      queue.push({ fn, resolve, reject });
      next();
    });
}

// the unbounded anti-pattern this replaces:
// await Promise.all(urls.map(fetch))  // fires ALL requests at once

Plain Promise.all(map) is the anti-pattern — nothing throttles it. In the libraries the knob is spelled `concurrency: n` (p-map, bluebird's Promise.map); this hand-rolled limiter is what p-limit does, ~15 lines. finally() is the release path, so a rejection still frees its slot.

TSTypeScript
export async function pool<T>(
  tasks: (() => Promise<T>)[],
  limit: number,
): Promise<T[]> {
  const results = new Array<T>(tasks.length);
  let next = 0;

  async function worker(): Promise<void> {
    while (next < tasks.length) {
      const i = next++; // claim BEFORE awaiting — the whole trick
      results[i] = await tasks[i]();
    }
  }

  const n = Math.min(Math.max(limit, 1), tasks.length);
  await Promise.all(Array.from({ length: n }, worker));
  return results; // input order preserved, unlike completion order
}

N workers pull indices off a shared cursor — no batches, so a straggler never stalls a lane. The Promise.all inside rejects the pool on the FIRST error: surviving workers finish their current task (their promises are handled, no unhandled-rejection noise) but no new task starts. Wrap each task in a result-catcher if you want gather-style all-outcomes.

GoGo
sem := make(chan struct{}, n) // acquire = send, release = receive
var wg sync.WaitGroup

for _, task := range tasks {
    wg.Add(1)
    go func(t Task) {
        defer wg.Done()
        sem <- struct{}{}        // blocks until a slot frees
        defer func() { <-sem }() // release even on panic
        _ = process(t)           // first error: see the note
    }(task)
}
wg.Wait()

A buffered channel of struct{} IS a semaphore — size N, send acquires, receive releases (defer both directions). The production one-liner is errgroup.SetLimit(n) with g.Go(...): same limiting plus first-error propagation and context cancellation that this shape lacks — reach for it outside teaching examples.

RsRust
use futures::stream::{self, StreamExt};

let results: Vec<Item> = stream::iter(tasks)
    .map(|t| async move { run(t).await })
    .buffer_unordered(n) // at most n futures in flight, started lazily
    .try_collect()       // Result<Vec<Item>, Error> — first error ends it
    .await?;

buffer_unordered LOSES ORDER — items come out in completion order, not input order; use .buffered(n) when positions must match the input. try_collect propagates the first error and drops the stream. The semaphore shape also exists: tokio::sync::Semaphore with owned permits (acquire_owned + a scope guard to release).

PHPPHP
// No stdlib async: parallelism means processes. pcntl_fork gives N
// workers, but fork COPIES the queue — pre-shard the work instead:
$shards = array_chunk($jobs, (int) ceil(count($jobs) / $n));

foreach ($shards as $shard) {
    $pid = pcntl_fork();
    if ($pid === 0) {               // child: run one shard, then exit
        foreach ($shard as $job) {
            run($job);
        }
        exit(0);
    }
}

while (pcntl_waitpid(0, $status) !== -1); // reap every child

The ecosystem gap is real: there is no portable stdlib answer. True shared-queue pools need the parallel extension or an event loop (ReactPHP, AMPHP, Swoole) for I/O-bound work — array_chunk sharding is the honest stdlib core. pcntl is Unix-only; the waitpid loop must reap children or they linger as zombies.

PyPython
import asyncio

async def run_limited(coros, n):
    sem = asyncio.Semaphore(n)

    async def guarded(coro):
        async with sem:            # acquire before, release after
            return await coro

    return await asyncio.gather(*map(guarded, coros))

# threads for blocking work — the cap IS the whole pool:
from concurrent.futures import ThreadPoolExecutor

with ThreadPoolExecutor(max_workers=n) as ex:
    results = list(ex.map(fn, items))  # ordered, all-or-wait

gather returns everything at once; as_completed(futures) is the streaming form — yields each result the moment it finishes. The semaphore gates EXECUTION, not task creation: every guarded() starts immediately and parks on acquire, which is fine until the task list itself is huge.

C#C#
using System.Threading;

// modern one-call answer (.NET 6+):
await Parallel.ForEachAsync(
    urls,
    new ParallelOptions { MaxDegreeOfParallelism = n },
    async (url, ct) => { await FetchAsync(url, ct); });

// the manual semaphore shape, when you need the values back:
using var sem = new SemaphoreSlim(n);
var results = await Task.WhenAll(urls.Select(async url =>
{
    await sem.WaitAsync();
    try     { return await FetchAsync(url); }
    finally { sem.Release(); } // the ONLY correct release site
}));

Parallel.ForEachAsync is the modern answer — MaxDegreeOfParallelism caps in-flight work and the await completes when all iterations do. In the manual shape, Release must sit in finally or one failure permanently shrinks the pool; Task.WhenAll propagates the first exception.

JvJava
try (ExecutorService pool = Executors.newFixedThreadPool(n)) {
    List<Callable<Item>> callables = tasks.stream()
            .map(t -> () -> process(t))
            .toList();

    List<Future<Item>> futures = pool.invokeAll(callables); // waits for ALL
    List<Item> results = new ArrayList<>();
    for (Future<Item> f : futures) {
        results.add(f.get()); // ExecutionException wraps the task's throw
    }
}

invokeAll BLOCKS until every task completes — no streaming. For results as they finish, ExecutorCompletionService: completionService.take().get() yields each in completion order. f.get() surfaces the first failure in ITERATION order, not completion order — cancel what's left yourself if that matters.

SwSwift
func runLimited<T: Sendable>(
    jobs: [@Sendable () async throws -> T],
    limit: Int
) async throws -> [T] {
    var results = [T?](repeating: nil, count: jobs.count)

    try await withThrowingTaskGroup(of: (Int, T).self) { group in
        var next = 0, inFlight = 0
        while next < jobs.count || inFlight > 0 {
            if next < jobs.count && inFlight < limit {
                let i = next
                group.addTask { (i, try await jobs[i]()) }
                next += 1
                inFlight += 1
            } else {
                // a slot must free before the next addTask
                if let (i, value) = try await group.next() {
                    results[i] = value
                    inFlight -= 1
                }
            }
        }
    }

    return results.compactMap { $0 }
}

TaskGroup has NO built-in limit — you count in-flight children yourself and only addTask when a slot is free; an AsyncSemaphore packages the same counting. group.next() yields children in completion order, so carry the index to keep results in input order. A throwing child cancels the group's remaining work.

KtKotlin
import kotlinx.coroutines.awaitAll
import kotlinx.coroutines.async
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.sync.Semaphore
import kotlinx.coroutines.sync.withPermit

suspend fun <T> runLimited(tasks: List<suspend () -> T>, n: Int): List<T> =
    coroutineScope {
        val sem = Semaphore(n)
        tasks.map { task ->
            async { sem.withPermit { task() } }
        }.awaitAll()
    }

Structured concurrency: the first failing async cancels its siblings and coroutineScope rethrows — no orphaned stragglers, the exact failure the naive chunk loop invites. awaitAll preserves input order. If you need every outcome despite failures, collect into a Result type inside each async.

RbRuby
queue = Queue.new
jobs.each { |job| queue << job }
n.times { queue << nil }   # one shutdown sentinel per worker

workers = Array.new(n) do
  Thread.new do
    while (job = queue.pop)    # Queue#pop BLOCKS when empty
      job.call
    end                        # nil sentinel ends the loop
  end
end
workers.each(&:join)

Queue#pop blocks forever on an empty queue — the nil sentinel is the shutdown protocol, one per worker, and missing it hangs join(). Under the GVL these threads overlap I/O waits, not CPU; processes (or Ractors) for parallel computation.

ZigZig
const std = @import("std");

const Pool = struct {
    mutex: std.Thread.Mutex = .{},
    cond: std.Thread.Condition = .{},
    jobs: std.ArrayList(Job), // the shared queue
    done: bool = false,

    fn worker(self: *Pool) void {
        while (true) {
            const job = blk: {
                self.mutex.lock();
                defer self.mutex.unlock();
                while (self.jobs.items.len == 0 and !self.done)
                    self.cond.wait(&self.mutex); // unlocks while waiting
                if (self.jobs.items.len == 0) return; // done + drained
                break :blk self.jobs.orderedRemove(0);
            };
            job.run(); // OUTSIDE the lock — keep critical sections tiny
        }
    }
};

// spawn n workers:  for (0..n) |i| threads[i] = try std.Thread.spawn(.{}, Pool.worker, .{&pool});
// shut down:        lock; done = true; cond.broadcast(); unlock; then join every thread.

Mutex + condition variable guard the queue: workers cond.wait while it is empty (the wait releases the lock and re-acquires on wake), pop one job under the lock, run it outside. Signal broadcast after setting done or sleeping workers never wake. std.Thread.Condition and std.Thread.spawn are the stdlib primitives — channels are not.