Ograniczanie współbieżności w Go: errgroup, puli workerów i backpressure

Ograniczaj pracę w toku, a nie tylko ją anuluj.

Page content

Nielimitowane rozgałęzienie jest domyślnym trybem awarii współbieżności w Go: każde wyrażenie go to kolejne zadanie w toku, bez limitu i bez właściciela. Ten artykuł dotyczy nadania im obu.

Anulowanie i ograniczanie to oddzielne mechanizmy sterowania. Anulowanie decyduje, kiedy praca się kończy; ograniczanie decyduje, ile z niej istnieje w danym momencie. Usługa może działać bezbłędnie pod kątem anulowania i nadal zawieść, ponieważ dziesięć tysięcy gorut, niezależnie od siebie, zadecydowało, że właśnie teraz jest odpowiedni moment, aby wywołać bazę danych.

szklany lejek przepuszczający strumień świecących cząstek przez regulowaną bramę

Strony powiązane z tą dotyczą drugiej połowy problemu. Go context.Context Done Right traktuje kontekst jako płaszczyznę sterowania zatrzymywaniem pracy; tutaj przedmiotem jest płaszczyzna pojemności — errgroup.SetLimit, kanały semaforowe, pula pracujących i ograniczone kolejki — oraz jeden kształt wycieku, który przetrwa errgroup.Wait, nawet gdy wszystkie ścieżki anulowania wydają się być prawidłowo podłączone.

Pięć prymitywów, jedna tabela decyzyjna

Prymityw Częściowe wyniki Awaria przy pierwszej pomyłce Backpressure (odpływ) Dynamiczna praca Najlepsze zastosowanie
sync.WaitGroup + errors.Join tak nie nie nie partie oczekiwania na wszystkie, w których każda wartość ma znaczenie
errgroup.WithContext nie tak nie nie rozgałęzienie, w którym pierwszy błąd kończy żądanie
errgroup.SetLimit(n) nie tak tak (blokuje Go) nie ograniczone rozgałęzienie z awarią przy pierwszej pomyłce
Kanał semaforowy tak ręcznie tak (blokuje akwizycję) tak niestandardowe sterowanie dopuszczeniem
Pula pracujących (Worker pool) tak ręcznie tak (pełna kolejka) tak stabilne strumienie, długotrwałe pracujące

Dwa pytania pozwalają wybrać wiersz. Po pierwsze: gdy jedno zadanie się nie powiedzie, czy nadal chcesz mieć wyniki pozostałych? Jeśli tak, semantyka grupy z awarią przy pierwszej pomyłce odrzuci pracę, której potrzebujesz. Po drugie: czy praca przychodzi ciągle, czy jest to ustalona partia, którą można policzyć przed uruchomieniem? Grupy obsługują partie; pula i kolejki obsługują strumienie.

errors.Join to element, który czyni wiersz z zwykłym WaitGroup wykonalnym — zbiera ono błędy każdego pracującego, zamiast zachowywać tylko pierwszy, co łączy się z regułami tłumaczenia granic opisanymi w Go Error Handling Architecture.

Wyciek, który przetrwa errgroup.Wait

Najlepiej znany tryb awarii errgroup nie ma nic wspólnego z ograniczaniem. To luka w anulowaniu, która ukrywa się za kodem, który wygląda na poprawny:

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

go func() { // producent — nikt nie jest właścicielem tego goruta
    for i := 0; i < 100; i++ {
        results <- i
    }
    close(results)
}()

g.Go(func() error { // konsument — członek grupy
    for r := range results {
        if r == 5 {
            return fmt.Errorf("save failed")
        }
    }
    return nil
})

err := g.Wait() // zwraca przy błędzie konsumenta

Gdy konsument zwraca, gctx jest anulowane, a Wait zwraca. Producent jest wstrzymany na results <- i — a anulowanie kontekstu nie odblokowuje wysyłki na kanale. Ponieważ nikt nie jest właścicielem tego goruta, pozostaje ono wstrzymane przez całe życie procesu. Uruchom ten wzór w pętli, a wyciek będzie dokładnie liniowy:

Wariant Aktywne goruty po 500 iteracjach
Wysyłka bez ramienia ctx.Done 501 (+500 wyciekłych, po jednym na iterację)
Wysyłka zawinięta w select z <-gctx.Done() 502 (+1 bazowe)
// naprawa: każda operacja na kanale wewnątrz zadania, które może zostać anulowane,
// otrzymuje ramię Done
select {
case results <- i:
case <-gctx.Done():
    return
}

Ta sama luka istnieje po stronie odbierania. Reguła jest mechaniczna: każda operacja na kanale w kodzie, który może zostać anulowany, otrzymuje ramię <-ctx.Done(), a każdy gorut ma właściciela, który na niego czeka lub go anuluje. goleak w CI łapie przypadki, których przegląd kodu przegapi — pasuje do tego samego automatycznego bramy jakości jak narzędzia z Go Linters: Essential Tools for Code Quality.

Warto nazwać dwa powiązane wzorce. Jeśli producent jest członkiem grupy, ten sam kod powoduje nie wyciek, ale deadlock — Wait blokuje się na zawsze, czekając na wstrzymaną wysyłkę, co co najmniej głośno sygnalizuje błąd. A producent, który wysyła za pomocą select, ale odbiera z kanału, którego konsument zakończył pracę, wycieka w ten sam sposób; ramiona Done są potrzebne po obu stronach.

errgroup.SetLimit: ograniczone rozgałęzienie z awarią przy pierwszej pomyłce

SetLimit zamienia errgroup w grupę ze sterowaniem dopuszczeniem. Każde wywołanie Go przekraczające limit blokuje się, aż zwolni się miejsce:

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

Trzy znaczenia są ważne w produkcji. Go blokuje się, więc pętla powyżej stosuje backpressure na producencie — to zazwyczaj jest to, czego chcesz, ale oznacza to, że gorut wywołującego może się zatrzymać, więc obsługa żądania zasilaająca dużą partię potrzebuje własnego budżetu timeoutu wokół całej pętli. SetLimit(0) blokuje każde wywołanie Go na zawsze. A zmiana limitu, gdy goruty są aktywne, powoduje panic — limit jest ustalony na całe życie grupy.

SetLimit dziedziczy semantykę grupy: pierwszy błąd anuluje kontekst, a Wait zwraca pierwszy błąd, odrzucając wyniki pracy wciąż w toku. Jest to właściwy prymityw, gdy awaria całej operacji przy pierwszej pomyłce jest poprawna — pobieranie jednej strony danych, wywoływanie zestawu niezależnych sprawdzeń stanu, rozgałęzianie żądania do shardów.

Kanały semaforowe: backpressure jako funkcja

Buforowany kanał o pojemności n to semafor bez zależności:

sem := make(chan struct{}, 8)
var wg sync.WaitGroup
for _, item := range items {
    wg.Add(1)
    sem <- struct{}{} // blokuje się, gdy 8 zadań jest w toku
    go func() {
        defer wg.Done()
        defer func() { <-sem }()
        _ = process(ctx, item)
    }()
}
wg.Wait()

Akwinizycja (zajęcie) odbywa się przed wyrażeniem go, więc goruty są tworzony wyłącznie wtedy, gdy istnieje miejsce. Rozmiar bufora jest umową: jest jednocześnie górnym limitem współbieżności i kolejką oczekujących na dopuszczenie.

Prymityw zasługuje na swoje miejsce, gdy potrzebujesz sterowania dopuszczeniem, którego grupy nie wyrażają — akwinizyjni z uwzględnieniem kontekstu, pula na każdego odbiorcę lub częściowe wyniki przy awarii:

select {
case sem <- struct{}{}: // miejsce zajęte
case <-ctx.Done():      // wywołujący zrezygnował przed wejściem
    return ctx.Err()
}

Zmierzone na tym samym obciążeniu, kanał semaforowy i SetLimit(8) kończą w różnicy poniżej 1,5% — wybór między nimi to kwestia semantyki, nie wydajności.

Pule pracujących (Worker pools): do strumieni, nie do partii

Stała pula długotrwałych pracujących oddziela dopuszczenie (kolejka) od wykonania (pracujący):

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
}

Pule pasują do pracy, która przychodzi ciągle — konsumenci kolejek, pętle pytań, pracujący żądań — tam, gdzie grupa na partię byłaby tworzona i likwidowana w nieskończoność. Od Go 1.25, sync.WaitGroup.Go usuwa szablonowy kod parowania Add/Done dla typowego kształtu. Umową shutdownu jest ta część, którą należy dobrze ogarnąć: zamknięcie jobs kończy pracujących po opróżnieniu kolejki; anulowanie ctx porzuca pracę w kolejce. Obie są uzasadnione, ale tylko jedna może być udokumentowanym zachowaniem.

Ograniczone kolejki: blokować lub odrzucać

Buforowany kanał jest też kolejką z sufitem. Gdy bufor się wypełni, producent blokuje się — a ten moment to decyzja projektowa. Blokowanie propaguje presję w górę do tego, kto zasila kolejkę; odrzucanie zmniejsza obciążenie kosztem utraconej pracy. Rurociąg monitoringu może odrzucać; rurociąg zamówień musi blokować lub przelewać do trwałego magazynu.

select {
case queue <- event: // dopuszczony
default:            // kolejka pełna: polityka odrzuceń znajduje się tutaj
    dropped.Add(1)
}

Niezależnie od polityki, uczyn ją jawną i mierzalną. Ciche ramię default to miejsce, gdzie zdarzenia znikają.

Czego kosztuje ograniczanie — zmierzone

To samo obciążenie, 10 000 zadań po 5 ms symulowanego I/O każde, Go 1.27.1:

Strategia Szczytowa liczba gorutów Czas rzeczywisty
Nielimitowane go na zadanie ~10 000 uruchomionych 18 ms
errgroup.SetLimit(8) 8 6,9 s
Kanał semaforowy (cap 8) 8 6,6 s

Wykonanie nielimitowane wygrywa czasem rzeczywistym, ponieważ nic w eksperymencie nie reaguje — dziesięć tysięcy gorutów timerów po prostu śpi równolegle. To jest dokładnie pułapka. Kosztem nielimitowanego rozgałęzienia nigdy nie jest CPU w Twoim własnym procesie; to dziesięć tysięcy jednoczesnych połączeń, odbiorca, który zaczyna tracić czas na timeouty pod wpływem skoku, i wykres pamięci, który podąża za głębokością kolejki tego, kogo wywołujesz. Na realnej zależności wykonanie nielimitowane nie kończy się w 18 ms — kończy timeoutem. Wykonania ograniczone płacą 6,6 sekund, ponieważ 10 000 zadań ÷ 8 slotów × 5 ms to 6,25 s nieuniknionej serializacji, a limit jest tu kluczem: zamienia niekontrolowany skok w przewidywalne opróżnianie 6,25-sekundowe.

Jedno doskonalenie ma znaczenie, gdy limit jest wyższy niż kilka slotów: utrzymuj limit na poziomie lub poniżej tego, co odbiorca może utrzymać współbieżnie, i wyprowadzaj go z tego limitu (rozmiar puli połączeń, limit przepustowości, pojemność pracujących), a nie z uczucia o „rozsądnej” równoległości.

Mierzenie i testowanie kodu ograniczonego

Trzy sygnały pokrywają większość przypadków. Trendy liczby gorutów (/sched/goroutines:goroutines eksportowane przez zbieracza Go w Prometheus) łapą wycieki jako nachylenie. Opóźnienia harmonizatora (/sched/latencies:seconds) łapią nasycenie — goruty gotowe do pracy czekające na czas CPU to realna presja niezależnie od liczby. Delta pprof dwóch zrzutów gorutów wskazuje precyzyjnie, gdzie mieszkają wstrzymane goruty.

W testach, ograniczeni pracujący to przypadek, do którego powstał Testing Concurrent Go Code with synctest: fałszywy czas sprawia, że pełne opróżnianie 10 000 zadań wykonuje się w milisekundach, synctest.Wait zastępuje sen i nadzieję na ciszę, a bąbel szybko się zawiesza, gdy pracujący pozostaje trwale zablokowany — powtórka wycieku powyżej umiera w testowym bąblu zamiast w produkcji.

W skali usługi to samo pytanie o ograniczanie pojawia się ponownie między usługami, gdzie semafor staje się limitem przepustowości, a wyłącznik awaryjny (circuit breaker) chroni odbiorcę — Circuit Breaker Pattern in Go opisuje tę warstwę, a Go Microservices for AI/ML Orchestration pokazuje formy oparte na kolejkach, jakie przybiera warstwa orkiestracji.

Wybór

flowchart TD A[Praca do wykonania współbieżnie] --> B{Stała partia czy ciągły strumień?} B -- ciągły strumień --> C[Pula pracujących lub ograniczona kolejka] B -- stała partia --> D{Wymagane częściowe wyniki
przy pierwszej awarii?} D -- tak --> E[WaitGroup + errors.Join
lub kanał semaforowy] D -- nie --> F[errgroup.WithContext] F --> G{Wymagany górny limit współbieżności?} G -- tak --> H[errgroup.SetLimit] G -- nie --> I[Grupa bez limitu
tylko gdy partia jest udowodnionym mała] C --> J{Pełna kolejka: blokować czy odrzucać?} J -- blokować --> K[Backpressure w górę] J -- odrzucać --> L[Zmniejsz obciążenie, licz odrzucenia]

Tabela i diagram sprowadzają się do jednego nawyku: przed napisaniem go, nazwij limit i właściciela. Limitem jest SetLimit, semafor lub ograniczona kolejka; właścicielem jest Wait, ramię ctx.Done na każdej operacji kanałowej, lub oba. Każde miejsce tworzenia goruta, które potrafi odpowiedzieć na te dwa nazwy w jednym rzucie oka, to takie, które nie wybudzi nikogo o 3 rano.

Przydatne linki

Subskrybuj

Otrzymuj nowe wpisy o systemach, infrastrukturze i inżynierii AI.