Limiter la concurrence en Go : errgroup, pools de travailleurs et backpressure

Limitez les travaux en cours, ne se contentez pas de les annuler.

Sommaire

Le fan-out non borné est le mode de défaillance par défaut de la concurrence en Go : chaque instruction go ajoute une tâche supplémentaire en cours, sans limite ni responsable. Cet article aborde la manière d’apporter ces deux éléments.

L’annulation et la limitation sont des contrôles distincts. L’annulation décide à quel moment le travail s’arrête ; la limitation décide de la quantité de travail existant à la fois. Un service peut annuler parfaitement et pourtant s’effondrer parce que dix mille goroutines ont chacune décidé, indépendamment, que c’était le bon moment pour appeler la base de données.

entonnoir de verre régulant un flux de particules lumineuses à travers une porte régulée

Les pages autour de celle-ci couvrent l’autre moitié du problème. Go context.Context Bien Fait traite le contexte comme le plan de contrôle pour l’arrêt du travail ; ici, le sujet est le plan de capacité — errgroup.SetLimit, les canaux de sémaphore, les pools de travailleurs et les files d’attente bornées — ainsi que la seule forme de fuite qui survit à errgroup.Wait même lorsque tous les chemins d’annulation semblent câblés.

Cinq primitives, une table de décision

Primitive Résultats partiels Échec rapide Backpressure Travail dynamique Idéal pour
sync.WaitGroup + errors.Join oui non non non lots attendus où chaque résultat compte
errgroup.WithContext non oui non non fan-out où la première erreur termine la requête
errgroup.SetLimit(n) non oui oui (bloquant Go) non fan-out borné avec échec rapide
Canal sémaphore oui manuel oui (bloque l’acquisition) oui contrôle d’admission personnalisé
Pool de travailleurs oui manuel oui (file pleine) oui flux réguliers, travailleurs à longue vie

Deux questions permettent de choisir la ligne. Première : lorsqu’une tâche échoue, souhaitez-vous toujours les résultats des autres ? Si oui, la sémantique de groupe à échec rapide supprime le travail dont vous aviez besoin. Seconde : le travail arrive-t-il en continu, ou s’agit-il d’un lot fixe que vous pouvez compter avant de lancer ? Les groupes gèrent les lots ; les pools et les files gèrent les flux.

errors.Join est ce qui rend la ligne WaitGroup simple viable — il collecte l’erreur de chaque travailleur au lieu de ne conserver que la première, ce qui s’associe avec les règles de traduction de frontières dans Architecture de Gestion des Erreurs en Go.

La fuite qui survit à errgroup.Wait

Le mode de défaillance d’errgroup le mieux connu n’a rien à voir avec la limitation. C’est un vide d’annulation, et il se cache derrière un code qui semble correct :

g, gctx := errgroup.WithContext(ctx)
results := make(chan int) // non tamponné

go func() { // producteur — personne ne possède cette goroutine
    for i := 0; i < 100; i++ {
        results <- i
    }
    close(results)
}()

g.Go(func() error { // consommateur — membre du groupe
    for r := range results {
        if r == 5 {
            return fmt.Errorf("save failed")
        }
    }
    return nil
})

err := g.Wait() // retourne à l'erreur du consommateur

Lorsque le consommateur retourne, gctx est annulé et Wait retourne. Le producteur est stationné sur results <- i — et l’annulation du contexte ne débloque pas l’envoi sur un canal. Personne ne possède cette goroutine, donc elle reste stationnée pour la vie du processus. Exécutez cette forme en boucle et la fuite est exactement linéaire :

Variante Goroutines actives après 500 itérations
Envoi sans une branche ctx.Done 501 (+500 fuites, une par itération)
Envoi enveloppé dans un select avec <-gctx.Done() 502 (+1 de base)
// le correctif : chaque opération de canal dans une tâche annulable
// obtient une branche Done
select {
case results <- i:
case <-gctx.Done():
    return
}

Le même vide existe côté réception. La règle est mécanique : chaque opération de canal dans du code qui peut être annulé obtient une branche <-ctx.Done(), et chaque goroutine a un propriétaire qui attend dessus ou l’annule. goleak dans CI capture les cas que la revue manque — il s’intègre dans la même porte de qualité automatisée que les outils dans Linters Go : Outils Essentiels pour la Qualité du Code.

Deux formes associées méritent d’être nommées. Si le producteur est un membre du groupe, le même code produit non pas une fuite mais un blocage — Wait bloque indéfiniment en attendant l’envoi stationné, ce qui échoue au moins bruyamment. Et un producteur qui envoie avec select mais reçoit sur un canal dont le consommateur a quitté fuit de la même manière ; les branches Done appartiennent aux deux côtés.

errgroup.SetLimit : fan-out borné avec échec rapide

SetLimit transforme un errgroup en groupe à contrôle d’admission. Chaque appel Go au-delà de la limite bloque jusqu’à ce qu’un emplacement se libère :

g, ctx := errgroup.WithContext(ctx)
g.SetLimit(8)
for _, item := range items {
    g.Go(func() error {
        return process(ctx, item)
    })
}
err := g.Wait()

Trois sémantiques importent en production. Go bloque, donc la boucle ci-dessus applique un backpressure au producteur — c’est généralement ce que vous voulez, mais cela signifie que la goroutine de l’appelant peut se bloquer, donc un gestionnaire de requête alimentant un grand lot a besoin de son propre budget de délai autour de la boucle complète. SetLimit(0) bloque chaque appel Go indéfiniment. Et changer la limite tandis que des goroutines sont actives provoque une panique — la limite est fixe pour la vie du groupe.

SetLimit hérite de la sémantique de groupe : la première erreur annule le contexte et Wait retourne la première erreur, supprimant les résultats du travail encore en cours. C’est la primitive appropriée lorsque l’échec de l’opération complète à la première défaillance est correct — récupérer une page de données, appeler un ensemble de vérifications de santé indépendantes, fan-out une requête vers des shards.

Canaux sémaphore : backpressure comme fonctionnalité

Un canal tamponné de capacité n est un sémaphore sans dépendances :

sem := make(chan struct{}, 8)
var wg sync.WaitGroup
for _, item := range items {
    wg.Add(1)
    sem <- struct{}{} // bloque lorsque 8 tâches sont en cours
    go func() {
        defer wg.Done()
        defer func() { <-sem }()
        _ = process(ctx, item)
    }()
}
wg.Wait()

L’acquisition a lieu avant l’instruction go, donc les goroutines ne sont créées que lorsqu’un emplacement existe. La taille du tampon est le contrat : c’est simultanément le plafond de concurrence et la file des admissions en attente.

La primitive gagne sa place lorsque vous avez besoin d’un contrôle d’admission que les groupes n’expriment pas — acquisition consciente du contexte, pools par flux descendant, ou résultats partiels en cas d’échec :

select {
case sem <- struct{}{}: // emplacement acquis
case <-ctx.Done():      // l'appelant a renoncé avant d'entrer
    return ctx.Err()
}

Mesuré sur le même volume de travail, le canal sémaphore et SetLimit(8) terminent à moins de 1,5 % l’un de l’autre — le choix entre eux est de sémantique, non de performance.

Pools de travailleurs : pour les flux, non pour les lots

Un pool fixe de travailleurs à longue vie sépare l’admission (la file) de l’exécution (les travailleurs) :

func pool(ctx context.Context, workers int) chan<- func() {
    jobs := make(chan func())
    var wg sync.WaitGroup
    for i := 0; i < workers; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            for {
                select {
                case f, ok := <-jobs:
                    if !ok {
                        return
                    }
                    f()
                case <-ctx.Done():
                    return
                }
            }
        }()
    }
    return jobs
}

Les pools conviennent au travail qui arrive en continu — consommateurs de files, boucles de polling, travailleurs de requêtes — là où un groupe par lot serait créé et détruit indéfiniment. Depuis Go 1.25, sync.WaitGroup.Go supprime le boilerplate d’appariement Add/Done pour la forme courante. Le contrat d’arrêt est la partie à obtenir correctement : fermer jobs termine les travailleurs après le drainage de la file ; annuler ctx abandonne le travail en file. Les deux sont légitimes, mais seulement l’un peut être le comportement documenté.

Files d’attente bornées : bloquer ou abandonner

Un canal tamponné est aussi une file avec un plafond. Une fois le tampon rempli, le producteur bloque — et ce moment est la décision de conception. Le blocage propage la pression en amont à celui qui alimente la file ; l’abandon allège la charge au coût de travail perdu. Un pipeline de surveillance peut abandonner ; un pipeline de commandes doit bloquer ou déverser vers un stockage durable.

select {
case queue <- event: // admis
default:            // file pleine : la politique d'abandon vit ici
    dropped.Add(1)
}

Quelle que soit la politique, rendez-la explicite et mesurée. Une branche default silencieuse est l’endroit où les événements vont disparaître.

Ce que la limitation coûte — mesuré

Le même volume de travail, 10 000 tâches de 5 ms d’E/S simulées chacune, Go 1.27.1 :

Stratégie Pic de goroutines Temps mural
go non borné par tâche ~10 000 lancées 18 ms
errgroup.SetLimit(8) 8 6,9 s
Canal sémaphore (cap 8) 8 6,6 s

L’exécution non bornée gagne en temps mural parce que rien dans l’expérience ne fait de résistance — dix mille goroutines de minuteur dorment simplement en parallèle. C’est exactement le piège. Le coût du fan-out non borné n’est jamais le CPU dans votre propre processus ; c’est dix mille connexions simultanées, un flux descendant qui commence à dépasser les délais sous la rafale, et un graphe de mémoire qui suit la profondeur de file de celui que vous appelez. Sur une dépendance réelle, l’exécution non bornée ne termine pas en 18 ms — elle dépasse les délais. Les exécutions bornées paient 6,6 secondes parce que 10 000 tâches ÷ 8 emplacements × 5 ms font 6,25 s de sérialisation inévitable, et le plafond est le point clé : il transforme une rafale incontrôlée en un drainage prévisible de 6,25 secondes.

Une raffinement importe lorsque le plafond est supérieur à une poignée d’emplacements : gardez le plafond à ou en dessous de ce que le flux descendant peut soutenir en parallèle, et dérivez-le de cette limite (taille du pool de connexions, limite de débit, capacité des travailleurs) plutôt que d’un sentiment à propos du parallélisme « raisonnable ».

Mesurer et tester le code borné

Trois signaux couvrent la plupart de cela. Les tendances du nombre de goroutines (/sched/goroutines:goroutines exporté via le collecteur Prometheus Go) capturent les fuites comme une pente. La latence du planificateur (/sched/latencies:seconds) capture la saturation — des goroutines exécutables attendant du temps CPU est une pression réelle quel que soit le nombre. Un delta pprof de deux dumps de goroutines localise où vivent les goroutines stationnées.

Pour les tests, les travailleurs bornés sont le cas pour lequel Tester le Code Go Concurrent avec synctest a été construit : le faux temps fait exécuter un drainage complet de 10 000 tâches en millisecondes, synctest.Wait remplace le sommeil-et-espèce pour la quiescence, et une bulle échoue rapidement lorsqu’un travailleur reste durablement bloqué — la reproduction de la fuite ci-dessus meurt dans une bulle de test au lieu de la production.

À l’échelle du service, la même question de limitation réapparaît entre les services, où le sémaphore devient une limite de débit et le circuit breaker protège le flux descendant — Motif Circuit Breaker en Go couvre cette couche, et Microservices Go pour l’Orchestration IA/ML montre les formes à base de files que prend la couche d’orchestration.

Choisir

flowchart TD A[Travail à exécuter en concurrence] --> B{Lot fixe ou flux continu ?} B -- flux continu --> C[Pool de travailleurs ou file bornée] B -- lot fixe --> D{Résultats partiels nécessaires
à la première défaillance ?} D -- oui --> E[WaitGroup + errors.Join
ou canal sémaphore] D -- non --> F[errgroup.WithContext] F --> G{Besoin d'un plafond de concurrence ?} G -- oui --> H[errgroup.SetLimit] G -- non --> I[groupe sans limite
uniquement si le lot est prouvée petit] C --> J{File pleine : bloquer ou abandonner ?} J -- bloquer --> K[Backpressure en amont] J -- abandonner --> L[Alléger la charge, compter les abandons]

La table et le diagramme de flux se compriment en une habitude : avant d’écrire go, nommer le plafond et le propriétaire. Le plafond est SetLimit, un sémaphore ou une file bornée ; le propriétaire est un Wait, une branche ctx.Done sur chaque opération de canal, ou les deux. Chaque site de création de goroutine qui peut répondre à ces deux noms d’un coup d’œil est celui qui ne rappellera personne à 3 h du matin.

Liens utiles

S'abonner

Recevez de nouveaux articles sur les systèmes, l'infrastructure et l'ingénierie IA.