Le Worker Pool est un pattern classique de la programmation concurrente : un groupe de goroutines (les workers) qui consomment des tâches depuis une file d’attente partagée.
En Go, ce pattern s’implémente en quelques lignes grâce aux channels et aux goroutines. C’est l’un des premiers réflexes à avoir quand on doit traiter un volume de tâches en parallèle sans exploser les ressources.

🎯 Quand utiliser un Worker Pool ?

  • Traitement de fichiers en masse — limiter les accès disque simultanés
  • Appels HTTP vers une API externe — respecter les rate limits, éviter de saturer la connexion
  • Jobs de fond (emails, notifications) — contrôler la charge sur les services downstream
  • Import/export de données — paralléliser sans OOM ni surcharge CPU

Le principe : borner le parallélisme pour éviter d’exploser la mémoire ou de saturer une ressource externe.

🛠️ L’implémentation minimale

Voici un Worker Pool fonctionnel en 30 lignes :

package main

import (
    "fmt"
    "sync"
)

func worker(id int, jobs <-chan int, results chan<- int, wg *sync.WaitGroup) {
    defer wg.Done()
    for job := range jobs {
        results <- job * 2 // traitement : doubler la valeur
        fmt.Printf("worker %d processed job %d\n", id, job)
    }
}

func main() {
    const numWorkers = 3
    const numJobs = 10

    jobs := make(chan int, numJobs)
    results := make(chan int, numJobs)
    var wg sync.WaitGroup

    // Lancer les workers
    for i := 1; i <= numWorkers; i++ {
        wg.Add(1)
        go worker(i, jobs, results, &wg)
    }

    // Envoyer les jobs
    for j := 1; j <= numJobs; j++ {
        jobs <- j
    }
    close(jobs) // signale aux workers qu'il n'y a plus de jobs

    // Attendre la fin des workers, puis fermer results
    go func() {
        wg.Wait()
        close(results)
    }()

    // Collecter les résultats
    for result := range results {
        fmt.Println("result:", result)
    }
}

Ce qui se passe :

  • On crée un channel jobs pour distribuer le travail
  • On lance numWorkers goroutines qui consomment depuis jobs
  • Chaque worker traite les jobs jusqu’à ce que le channel soit fermé
  • Le WaitGroup synchronise la fin de tous les workers
  • On ferme results quand tous les workers ont terminé

On voit donc que Go fournit tout ce qu’il faut sans avoir besoin de bibliothèque externe ou d’abstraction complexe.

ℹ️ A NOTER

En lançant cet exemple plusieurs fois, l’ordre des logs change à chaque run, et un seul worker semble parfois traiter presque tous les jobs pendant que les autres restent inactifs.

Ce n’est pas un bug 🐛

👉 jobs et results sont bufferisés à la taille exacte du nombre de jobs, donc le programme principal empile tous les jobs et ferme le channel avant même que les workers n’aient une chance de s’exécuter.

Le premier worker élu par le scheduler peut alors vider jobs sans jamais bloquer et sans point de blocage.
Go n’a aucune raison de changer de goroutine, donc ce worker rafle presque tout avant que les autres soient sollicités.

De la même façon, results <- job*2 et le fmt.Printf qui suit ne sont pas atomiques : rien n’empêche la fonction main de vider entièrement results avant que le worker n’ait eu le temps d’afficher son propre log, qui peut donc apparaître après coup.

Pour observer un partage plus équilibré entre workers, il suffit de retirer le buffer des channels (make(chan int)) ou d’ajouter un léger traitement (time.Sleep) dans le worker.

⚠️ Les pièges à éviter

1. Oublier de fermer le channel jobs

// ❌ Les workers bloquent indéfiniment
for j := range jobs {
    // ...
}
// Le range ne termine jamais si jobs n'est pas fermé !

Solution : Toujours appeler close(jobs) après avoir envoyé tous les jobs.

2. Capturer la variable de boucle (avant Go 1.22)

// ❌ Avant Go 1.22 : tous les workers reçoivent la même valeur !
for i := 1; i <= numWorkers; i++ {
    go func() {
        fmt.Println(i) // i vaut numWorkers+1 pour tous
    }()
}

Solution : Passer i en paramètre ou utiliser Go 1.22+ où ce problème est corrigé.

// ✅ Correct
for i := 1; i <= numWorkers; i++ {
    go func(id int) {
        fmt.Println(id)
    }(i)
}

3. Deadlock sur le channel results

Si vous écrivez dans results de manière synchrone avant de le lire :

// ❌ Deadlock si results est unbuffered et qu'on n'a pas de lecteur
results <- value // bloque si personne ne lit

Solution : Soit utiliser un buffer suffisant, soit lancer le lecteur dans une goroutine séparée avant d’écrire.

4. Ne pas propager le contexte

Dans du code de production, les workers doivent respecter l’annulation :

func worker(ctx context.Context, jobs <-chan Job) {
    for {
        select {
        case <-ctx.Done():
            return // arrêt propre
        case job, ok := <-jobs:
            if !ok {
                return
            }
            process(job)
        }
    }
}

Sans ça, vos workers tournent jusqu’à la fin du programme même si le contexte est annulé.

🚀 Vers une version plus robuste

Voici une version plus complète avec gestion du contexte et erreurs :

func RunWorkerPool[T any, R any](
    ctx context.Context,
    numWorkers int,
    jobs []T,
    process func(context.Context, T) (R, error),
) ([]R, error) {
    jobsCh := make(chan T)
    resultsCh := make(chan R, len(jobs))
    errCh := make(chan error, 1)

    var wg sync.WaitGroup

    // Lancer les workers
    for i := 0; i < numWorkers; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            for {
                select {
                case <-ctx.Done():
                    return
                case job, ok := <-jobsCh:
                    if !ok {
                        return
                    }
                    result, err := process(ctx, job)
                    if err != nil {
                        select {
                        case errCh <- err:
                        default:
                        }
                        return
                    }
                    resultsCh <- result
                }
            }
        }()
    }

    // Distribuer les jobs
    go func() {
        defer close(jobsCh)
        for _, job := range jobs {
            select {
            case <-ctx.Done():
                return
            case jobsCh <- job:
            }
        }
    }()

    // Attendre et fermer results
    go func() {
        wg.Wait()
        close(resultsCh)
    }()

    // Collecter
    var results []R
    for r := range resultsCh {
        results = append(results, r)
    }

    select {
    case err := <-errCh:
        return results, err
    default:
        return results, nil
    }
}

Cette version :

  • Respecte l’annulation via context.Context
  • Utilise les generics (Go 1.18+) pour être réutilisable
  • Gère la première erreur rencontrée
  • Arrête proprement les workers en cas d’annulation ou d’erreur

📝 En résumé

  • Simplicité — Channels + goroutines + WaitGroup = Worker Pool
  • Fermer les channels — Toujours close(jobs) après envoi
  • Contexte — Propager ctx pour l’annulation
  • Boucle variable — Passer en paramètre (ou Go 1.22+)
  • Buffer — Ajuster selon le cas d’usage

En Go, créer un worker pool ne requiert que la stdlib et quelques dizaines de lignes de code.

📚 Resources