Skip to content

Limiter la concurrence à N tâches à la fois snippet

Lancer mille tâches mais seulement N à la fois — le pool borné qui protège les services en aval et la mémoire.

Lancer mille tâches mais seulement N à la fois — le pool borné qui protège les services en aval et la mémoire. Le piège, c'est la boucle naïve par lots : await Promise.all(chunk) attend le PLUS LENT de chaque lot avant que le suivant ne démarre (un traînard bloque la voie), et un rejet peut laisser les autres orphelines. Les formes correctes : N workers qui tirent sur une file partagée, ou un sémaphore acquis avant et relâché après — et propage toujours la première erreur.

Recette exécutable · 12 langages
Concurrency & Parallelismconcurrencyasyncsemaphoreworker-poolparallelismrate-limiting

Every language

12 langages, 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.