Limitare la concorrenza in Go: errgroup, pool di worker e backpressure

Limita le operazioni in corso, non limitarti ad annullarle.

Indice

Il fan-out illimitato è la modalità di guasto predefinita nella concorrenza di Go: ogni istruzione go aggiunge un altro task in volo, senza un limite e senza un responsabile. Questo articolo si occupa di fornire entrambi gli aspetti.

L’annullamento e il limitazione sono controlli distinti. L’annullamento decide quando il lavoro si ferma; la limitazione decide quanta parte di esso esiste in un dato momento. Un servizio può annullare perfettamente il lavoro e comunque crollare perché diecimila goroutine hanno deciso, ciascuna in modo indipendente, che quel momento fosse buono per chiamare il database.

imbuto di vetro che limita il flusso di particelle luminose attraverso un varco regolato

Le pagine attorno a questa coprono l’altra metà del problema. Go context.Context fatto bene tratta il contesto come il piano di controllo per fermare il lavoro; qui l’argomento è il piano della capacità — errgroup.SetLimit, canali semaforo, pool di worker e code limitate — oltre alla forma di perdita che sopravvive a errgroup.Wait anche quando tutti i percorsi di annullamento sembrano correttamente collegati.

Cinque primitive, una tabella di scelta

Primitiva Risultati parziali Fail-fast Backpressure Lavoro dinamico Ideale per
sync.WaitGroup + errors.Join sì no no no batch di attesa in cui ogni risultato conta
errgroup.WithContext no sì no no fan-out in cui il primo errore termina la richiesta
errgroup.SetLimit(n) no sì sì (blocca Go) no fan-out fail-fast con limite
Canale semaforo sì manuale sì (blocca l’acquisizione) sì controllo di ammissione personalizzato
Pool di worker sì manuale sì (la code si riempie) sì flussi continui, worker a lunga durata

Due domande individuano la riga giusta. Prima: quando un task fallisce, si vogliono ancora i risultati degli altri? Se sì, la semantica fail-fast dei gruppi scarta lavoro di cui si aveva bisogno. Seconda: il lavoro arriva continuamente o è un batch fisso che si può contare prima di lanciarlo? I gruppi gestiscono i batch; pool e code gestiscono i flussi.

errors.Join è ciò che rende fattibile la riga del semplice WaitGroup — raccoglie l’errore di ogni worker invece di tenere solo il primo, il che si abbina alle regole di traduzione dei confini in Architettura della gestione degli errori in Go.

La perdita che sopravvive a errgroup.Wait

La modalità di guasto più nota di errgroup non ha nulla a che fare con la limitazione. È un buco nell’annullamento, e si nasconde dietro codice che sembra corretto:

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

go func() { // produttore — nessuno possiede questa goroutine
    for i := 0; i < 100; i++ {
        results <- i
    }
    close(results)
}()

g.Go(func() error { // consumatore — membro del gruppo
    for r := range results {
        if r == 5 {
            return fmt.Errorf("salvataggio fallito")
        }
    }
    return nil
})

err := g.Wait() // restituisce all'errore del consumatore

Quando il consumatore restituisce, gctx viene annullato e Wait restituisce. Il produttore è parcheggiato su results <- i — e l’annullamento del contesto non sblocca un invio su canale. Nessuno possiede quella goroutine, quindi resta parcheggiata per tutta la vita del processo. Se si esegue questa forma in un loop, la perdita è esattamente lineare:

Variante Goroutine attive dopo 500 iterazioni
Invio senza un ramo ctx.Done 501 (+500 perse, una per iterazione)
Invio avvolto in un select con <-gctx.Done() 502 (+1 di base)
// la correzione: ogni operazione su canale dentro un task annullabile
// ottiene un ramo Done
select {
case results <- i:
case <-gctx.Done():
    return
}

La stessa lacuna esiste lato ricezione. La regola è meccanica: ogni operazione su canale all’interno di codice che può essere annullato deve avere un ramo <-ctx.Done(), e ogni goroutine ha un responsabile che attende su di essa o la annulla. goleak in CI cattura i casi che la revisione del codice trascura — si inserisce nello stesso gate di qualità automatizzato degli strumenti in Go Linters: Strumenti essenziali per la qualità del codice.

Vale la pena nominare due forme correlate. Se il produttore è un membro del gruppo, lo stesso codice produce non una perdita ma un deadlock — Wait blocca in attesa dell’invio parcheggiato, il che almeno fallisce in modo evidente. E un produttore che invia con select ma riceve su un canale il cui consumatore si è già arrestato perde nello stesso modo; i rami Done devono essere presenti su entrambi i lati.

errgroup.SetLimit: fan-out fail-fast con limite

SetLimit trasforma un errgroup in un gruppo con controllo di ammissione. Ogni chiamata a Go oltre il limite blocca fino a quando uno slot non si libera:

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

Tre semantici contano in produzione. Go blocca, quindi il loop sopra applica il backpressure al produttore — è di solito ciò che si vuole, ma significa che la goroutine del chiamante può restare ferma, quindi un handler di richiesta che alimenta un batch grande ha bisogno di un proprio budget di timeout attorno a tutto il loop. SetLimit(0) blocca ogni chiamata a Go per sempre. E cambiare il limite mentre le goroutine sono attive provoca un panic — il limite è fisso per tutta la vita del gruppo.

SetLimit eredita la semantica del gruppo: il primo errore annulla il contesto e Wait restituisce il primo errore, scartando i risultati del lavoro ancora in volo. È la primitiva giusta quando fallire l’intera operazione al primo errore è corretto — recuperare una pagina di dati, chiamare un set di health check indipendenti,分发 una richiesta a sharded.

Canali semaforo: il backpressure come funzione

Un canale limitato di capacità n è un semaforo senza dipendenze:

sem := make(chan struct{}, 8)
var wg sync.WaitGroup
for _, item := range items {
    wg.Add(1)
    sem <- struct{}{} // blocca quando 8 task sono in volo
    go func() {
        defer wg.Done()
        defer func() { <-sem }()
        _ = process(ctx, item)
    }()
}
wg.Wait()

L’acquisizione avviene prima dell’istruzione go, quindi le goroutine vengono create solo quando esiste uno slot. La dimensione del buffer è il contratto: è simultaneamente il tetto della concorrenza e la code delle ammissioni in attesa.

La primitiva giustifica il suo posto quando serve un controllo di ammissione che i gruppi non esprimono — acquisizione consapevole del contesto, pool per downstream, o risultati parziali in caso di errore:

select {
case sem <- struct{}{}: // slot acquisito
case <-ctx.Done():      // il chiamante ha rinunciato prima di entrare
    return ctx.Err()
}

Misurato sullo stesso carico di lavoro, il canale semaforo e SetLimit(8) terminano entro il 1,5% l’uno dall’altro — la scelta tra loro è semantica, non performance.

Pool di worker: per flussi, non batch

Un pool fisso di worker a lunga durata separa l’ammissione (la code) dall’esecuzione (i worker):

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
}

I pool si adattano al lavoro che arriva continuamente — consumer di code, loop di polling, worker di richieste — dove un gruppo per batch verrebbe creato e smontato all’infinito. Da Go 1.25, sync.WaitGroup.Go rimuove la boilerplate di appaiamento Add/Done per la forma comune. Il contratto di spegnimento è la parte da gestire bene: chiudere jobs termina i worker dopo che la code si è svuotata; annullare ctx abbandona il lavoro in code. Entrambe sono legittime, ma solo una può essere il comportamento documentato.

Code limitate: bloccare o scartare

Un canale limitato è anche una code con un tetto. Una volta che il buffer si riempie, il produttore blocca — e quel momento è la decisione di design. Bloccare propaga la pressione a monte a chi alimenta la code; scartare allevia il carico a costo di lavoro perso. Un pipeline di monitoraggio può scartare; un pipeline di ordini deve bloccare o riversare su storage durevole.

select {
case queue <- event: // ammesso
default:            // code piena: la politica di scarto sta qui
    dropped.Add(1)
}

Qualunque sia la politica, rendila esplicita e misurata. Un ramo default silenzioso è dove gli eventi vanno a scomparire.

Cosa costa la limitazione — misurato

Lo stesso carico di lavoro, 10.000 task di 5 ms di I/O simulato ciascuno, Go 1.27.1:

Strategia Goroutine di picco Tempo reale
go illimitato per task ~10.000 avviate 18 ms
errgroup.SetLimit(8) 8 6,9 s
Canale semaforo (cap 8) 8 6,6 s

L’esecuzione illimitata vince sul tempo reale perché nulla nell’esperimento reagisce — diecimila goroutine di timer semplicemente dormono in parallelo. È esattamente la trappola. Il costo del fan-out illimitato non è mai la CPU nel proprio processo; sono diecimila connessioni simultanee, un downstream che inizia ad andare in timeout sotto il picco, e un grafico della memoria che segue la profondità della code di chi stai chiamando. Su una dipendenza reale, l’esecuzione illimitata non termina in 18 ms — va in timeout. Le esecuzioni limitate pagano 6,6 secondi perché 10.000 task ÷ 8 slot × 5 ms fa 6,25 s di serializzazione inevitabile, e il tetto è il punto: trasforma un picco non controllato in uno svuotamento prevedibile di 6,25 secondi.

Un raffinamento conta quando il tetto è più alto di un pugno di slot: mantieni il tetto a o sotto ciò che il downstream può sostenere in concomitanza, e derivalo da quel limite (dimensione del pool di connessioni, rate limit, capacità dei worker) piuttosto che da un’idea di “parallelismo ragionevole”.

Misurare e testare il codice limitato

Tre segnali coprono la maggior parte dei casi. Le tendenze del conteggio goroutine (/sched/goroutines:goroutines esportato via il collectore Go di Prometheus) catturano le perdite come una pendenza. La latenza del scheduler (/sched/latencies:seconds) cattura la saturazione — goroutine eseguibili in attesa di tempo CPU è pressione reale indipendentemente dal conteggio. Un delta pprof di due dump di goroutine individua dove vivono le goroutine parcheggiate.

Per i test, i worker limitati sono il caso per cui Testing Concurrent Go Code with synctest è stato costruito: il tempo finto fa sì che uno svuotamento completo di 10.000 task si esegua in millisecondi, synctest.Wait sostituisce lo “sleep e speranza” per la quiete, e una bolla fallisce rapidamente quando un worker resta durabilmente bloccato — la riproduzione della perdita sopra muore in una bolla di test invece che in produzione.

A livello di servizio, la stessa domanda di limitazione riappare tra i servizi, dove il semaforo diventa un rate limit e il circuit breaker protegge il downstream — Circuit Breaker Pattern in Go copre quel livello, e Go Microservices for AI/ML Orchestration mostra le forme basate su code che il tier di orchestrazione assume.

Scelta

flowchart TD A[Lavoro da eseguire in concomitanza] --> B[Batch fisso o flusso continuo?] B -- flusso continuo --> C[Pool di worker o code limitata] B -- batch fisso --> D[Risultati parziali necessari
al primo errore?] D -- sì --> E[WaitGroup + errors.Join
o canale semaforo] D -- no --> F[errgroup.WithContext] F --> G[Serve un tetto di concorrenza?] G -- sì --> H[errgroup.SetLimit] G -- no --> I[gruppo senza limite
solo quando il batch è provabilmente piccolo] C --> J[Code piena: bloccare o scartare?] J -- bloccare --> K[Backpressure a monte] J -- scartare --> L[Alleviare il carico, contare gli scarti]

La tabella e il diagramma a flusso si comprono in un’abitudine: prima di scrivere go, nominare il tetto e il responsabile. Il tetto è SetLimit, un semaforo, o una code limitata; il responsabile è un Wait, un ramo ctx.Done su ogni operazione di canale, o entrambi. Ogni sito di creazione di goroutine che può rispondere a quei due nomi in un solo sguardo è uno che non chiamerà nessuno alle 3 di notte.

Iscriviti

Ricevi nuovi articoli su sistemi, infrastruttura e ingegneria AI.