Перейти до змісту

Обмежена конкурентність

Горутини дешеві, тож очевидний цикл запускає одну горутину на кожен елемент. При тисячі елементів це тисяча горутин, які одночасно б'ють в одну й ту саму базу даних чи один і той самий API:

var wg sync.WaitGroup
for _, job := range jobs {
    wg.Add(1)
    go func() {
        defer wg.Done()
        process(job)
    }()
}
wg.Wait()

Для десяти елементів тут усе гаразд. А для десяти тисяч це спосіб вичерпати дескриптори файлів, вичерпати пул з'єднань або наштовхнутися на обмеження частоти запитів (rate limiting) — і збій приходить під навантаженням, а не під час тестування.

Буферизований канал — це семафор

Буферизований канал місткістю n пропускає n відправників, перш ніж заблокуватися. Це і є весь рахувальний семафор (counting semaphore). Захоплення — це відправлення, звільнення — отримання:

sem := make(chan struct{}, 3)        // at most 3 at once

var wg sync.WaitGroup
for _, job := range jobs {
    wg.Add(1)
    go func() {
        defer wg.Done()
        sem <- struct{}{}            // acquire: blocks when 3 are running
        defer func() { <-sem }()     // release
        process(job)
    }()
}
wg.Wait()

struct{} — це тип на нуль байтів зі статті про структури — канал не несе жодних даних, лише дозвіл. Запустіть цикл вище на тисячі завдань, підраховуючи, скільки їх одночасно всередині process, і пік буде рівно 3.

Усі горутини все одно створюються — вони просто чекають у черзі на захоплення. Це нормально: горутина, що чекає на каналі, тримає невеликий стек і не займає потік ОС — саме тому горутини й дешеві. Якщо навіть їх створення — це забагато, використовуйте натомість пул робітників, який запускає фіксовану кількість горутин і годує їх із каналу. Емпіричне правило:

Форма роботи Що використати
відомий обсяг роботи, хочете обмежити паралельність семафор
необмежений потік роботи пул робітників
робота надходить швидше, ніж ви встигаєте її обробляти пул робітників, щоб черга була видимою

Де розмістити захоплення

Захоплюйте всередині горутини, як вище, а не перед інструкцією go. Розміщення зовні блокує сам цикл, а це означає, що горутини створюються по одній, і все, що йде після циклу, теж чекає.

Звільняйте через defer, щоб паніка чи ранній return усередині process не могли протекти й втратити слот. Протеклий слот назавжди зменшує ліміт, а коли протече достатньо, усе входить у deadlock — баг, що виглядає як зависання, без жодної помилки де-небудь.

Збирання результатів

Запис у різні індекси заздалегідь розміреного зрізу взагалі не потребує блокування. Різні індекси — це різна пам'ять, тож гонки немає:

jobs := []string{"a", "b", "c", "d"}
results := make([]string, len(jobs))

sem := make(chan struct{}, 2)
var wg sync.WaitGroup
for i, job := range jobs {
    wg.Add(1)
    go func() {
        defer wg.Done()
        sem <- struct{}{}
        defer func() { <-sem }()
        results[i] = strings.ToUpper(job)
    }()
}
wg.Wait()
fmt.Println(results)   // output: [A B C D]

Результати лишаються в порядку вхідних даних, чого канал вам не дав би. Зауважте, що i та job — це змінні на кожну ітерацію, тож горутина захоплює саме потрібні значення — див. горутини.

А от усе спільне блокування таки потребує. Найпоширеніший випадок — зберегти першу помилку:

var mu sync.Mutex
var firstErr error

// inside each goroutine:
if _, err := process(job); err != nil {
    mu.Lock()
    if firstErr == nil {
        firstErr = err
    }
    mu.Unlock()
}

Цей патерн настільки поширений, що команда Go постачає для нього готовий помічник у golang.org/x/sync/errgroup. Це окремий модуль, а не частина стандартної бібліотеки, тож він розглядається разом з іншими сторонніми бібліотеками.

sync.Map

Звичайна мапа під захистом sync.Mutex — правильний варіант за замовчуванням. sync.Map — це окремий тип для двох конкретних випадків: ключ записується один раз, а читається багато разів, або різні горутини торкаються здебільшого непересічних ключів. Він міняє типи часу компіляції на any, тож сягайте по нього, лише коли профілювання каже, що варто.

var cache sync.Map

cache.Store("a", 1)
v, ok := cache.Load("a")
fmt.Println(v, ok)              // output: 1 true

actual, loaded := cache.LoadOrStore("b", 2)
fmt.Println(actual, loaded)     // output: 2 false

actual, loaded = cache.LoadOrStore("b", 99)
fmt.Println(actual, loaded)     // output: 2 true

cache.Delete("a")
_, ok = cache.Load("a")
fmt.Println(ok)                 // output: false

LoadOrStore — це той метод, що виправдовує своє існування: він атомарно повертає наявне значення, якщо воно є, і зберігає ваше, якщо його нема, тож дві горутини, що змагаються за заповнення одного ключа, погоджуються, хто переміг.

З погляду Python: у стандартній бібліотеці немає ThreadPoolExecutor(max_workers=3). Буферизований канал і є max_workers, а решту ви складаєте самі — тому ці самі шість рядків з'являються в кожній Go-кодовій базі.

Швидка довідка

Форма Значення
sem := make(chan struct{}, n) семафор, що дозволяє n одночасно
sem <- struct{}{} захоплення (блокує, коли повний)
defer func() { <-sem }() звільнення, безпечне при паніці
results[i] = v збирання результатів за індексом без блокувань
mu.Lock() навколо спільного стану усе, що не по індексу
sync.Map запис-раз-читання-багато або непересічні ключі

Джерела