Горутины
Основу параллелизма в Go составляют горутины — функции, запущенные с помощью ключевого слова go:
func main() {
var wg sync.WaitGroup
wg.Add(2)
go func() {
defer wg.Done()
fmt.Println("worker 1")
}()
go func() {
defer wg.Done()
fmt.Println("worker 2")
}()
wg.Wait()
}
worker 2
worker 1
Go runtime управляет этими горутинами и распределяет их между потоками операционной системы, работающими на ядрах процессора. По сравнению с потоками ОС, горутины легковесные, поэтому можно создавать сотни или тысячи их.
Горутины полностью независимы. Главная функция тоже является горутиной, но запускается неявно при старте программы. Когда main завершается, другие горутины также завершают работу.
В приведённом выше примере используется wait group (sync.WaitGroup) для ожидания завершения горутин. Wait group содержит внутри счётчик. Вызов Add(n) увеличивает его на n, а Done() уменьшает на один. Wait() блокирует вызывающую горутину (в данном случае main) до обнуления счётчика. Таким образом, main ждёт завершения обеих рабочих горутин перед выходом.
WaitGroup.Go автоматически увеличивает счётчик, запускает функцию в горутине и уменьшает счётчик после завершения:
func main() {
var wg sync.WaitGroup
wg.Go(func() {
fmt.Println("worker 1")
})
wg.Go(func() {
fmt.Println("worker 2")
})
wg.Wait()
}
worker 2
worker 1
Каналы
Горутины могут передавать значения друг другу через каналы. Канал — это как окно, через которое одна горутина может что-то бросить, а другая — поймать:
func main() {
messages := make(chan string)
go func() { messages <- "ping" }()
msg := <-messages
fmt.Println(msg)
}
ping
Отправка значения в канал — синхронная операция. Когда отправляющая горутина записывает значение в канал (ch <- val), она блокируется и ждёт, пока кто-нибудь получит это значение (<-ch). Только после этого она продолжает работу.
Выходной канал
Возвращение выходного канала из функции и заполнение его внутри внутренней горутины — распространённый паттерн в Go. Это позволяет вызывающей стороне получать значения через канал, в то время как владеющая функция сохраняет контроль над ним:
func generate(start, stop int) chan int {
out := make(chan int)
go func() {
for i := start; i < stop; i++ {
out <- i
}
}()
return out
}
Закрытие канала
Чтобы сигнализировать читателям о том, что все данные отправлены, отправляющая горутина закрывает канал с помощью close():
func generate(start, stop int) chan int {
out := make(chan int)
go func() {
defer close(out)
for i := start; i < stop; i++ {
out <- i
}
}()
return out
}
Читатель проверяет статус канала со вторым значением («запятая OK») при чтении:
func main() {
in := generate(5, 10)
for {
num, ok := <-in
if !ok {
break
}
fmt.Print(num, " ")
}
}
5 6 7 8 9
Пока канал открыт, читатель получает следующее значение и статус true. Если канал закрыт, читатель получает нулевое значение и статус false.
Канал можно закрыть только один раз. Повторное закрытие или запись в закрытый канал вызывает панику.
Единственная причина закрыть канал — сигнализировать его читателям о том, что все данные отправлены. Если это неважно для читателей, то закрывать канал не обязательно. Когда канал больше не используется, сборщик мусора Go освободит его ресурсы, закрыт он или нет.
Итерация по каналу
range автоматически читает следующее значение из канала и проверяет, закрыт ли он. Если канал закрыт, цикл завершается:
func main() {
nums := generate(5, 10)
for n := range nums {
fmt.Print(n, " ")
}
}
5 6 7 8 9
Range над каналом возвращает одно значение, а не пару, в отличие от range над срезом.
Направленные каналы
Можно защитить себя от случайных ошибок записи/закрытия, установив направление канала. Каналы могут быть:
chan(двунаправленный): для чтения и записи (по умолчанию);chan<-(только отправка): только для записи;<-chan(только приём): только для чтения.
Нельзя читать из канала «только отправка» или писать в канал «только приём» (и нельзя его закрывать).
Каналы обычно инициализируются для чтения и записи, а в параметрах функций указываются направлением. Go автоматически преобразует обычный канал в направленный:
stream := make(chan int)
go func(in chan<- int) {
in <- 42
}(stream)
func(out <-chan int) {
fmt.Println(<-out)
}(stream)
42
Буферизованные каналы
Буферизованные каналы работают как очередь FIFO с буфером фиксированного размера для хранения значений.
Пока буфер имеет свободное место, запись в канал не блокирует горутину. Аналогично, пока буфер содержит значения, чтение из канала не блокирует горутину:
stream := make(chan int, 3)
stream <- 11
stream <- 12
stream <- 13
fmt.Println(<-stream)
fmt.Println(<-stream)
11
12
По умолчанию, если размер буфера не указан, канал небуферизованный (размер буфера равен нулю).
Буферизованные каналы работают со встроенными функциями len() и cap():
stream := make(chan int, 3)
stream <- 11
fmt.Println(cap(stream), len(stream))
3 1
Чтение из закрытого буферизованного канала возвращает значения из буфера и статус true. После получения всех значений возвращается нулевое значение и статус false, как в обычном канале:
stream := make(chan int, 1)
stream <- 11
close(stream)
val, ok := <-stream
fmt.Println(val, ok)
// 11 true
val, ok = <-stream
fmt.Println(val, ok)
// 0 false
11 true
0 false
nil канал
Как и любой тип в Go, каналы имеют нулевое значение, которое равно nil.
Запись или чтение из nil канала блокирует горутину бесконечно:
var stream chan int
go func() {
// блокируется навсегда
stream <- 1
}()
// блокируется навсегда
<-stream
Закрытие nil канала вызывает панику:
var stream chan int
close(stream)
// panic: close of nil channel
Select
Select — это как switch, но специально разработан для каналов. Вот что он делает:
- Проверяет, какие ветки не заблокированы.
- Если несколько веток готовы, случайно выбирает одну для выполнения.
- Если все ветки заблокированы и есть ветка
default, выполняет её. - Если все ветки заблокированы и нет ветки
default, ждёт, пока одна из них будет готова.
Select используется для управления потоком данных в pipelines:
// merge отправляет значения из in1 и in2 в выходной канал.
func merge(in1, in2 <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for in1 != nil || in2 != nil {
select {
case val1, ok := <-in1:
if ok { out <- val1 } else { in1 = nil }
case val2, ok := <-in2:
if ok { out <- val2 } else { in2 = nil }
}
}
}()
return out
}
// Допустим, отправляем 10..12 в in1, 20..22 в in2,
// и вызываем merge(in1, in2)
10 11 20 12 21 22
Для отмены горутин:
// process изменяет значения из in и отправляет их в out
// до тех пор, пока in не исчерпается или cancel не закроется.
func process(cancel chan struct{}, in <-chan int) <-chan int {
out := make(chan int)
go func() {
for val := range in {
select {
case out <- val*10:
case <-cancel:
fmt.Println("canceled")
return
}
}
}()
return out
}
// Допустим, отправляем значения 11 и 12 в in
// и затем вызываем close(cancel)
110
120
canceled
Для неблокирующих операций:
// multiplier возвращает функцию, которая умножает
// входное значение на 10 и отправляет его в канал
// или возвращает ошибку, если канал занят.
func multiplier(ch chan<- int) func(n int) error {
return func(n int) error {
select {
case ch <- n*10:
return nil
default:
return errors.New("busy")
}
}
}
func main() {
nums := make(chan int, 1)
multiply := multiplier(nums)
err := multiply(11)
fmt.Println(<-nums, err)
// 110 <nil>
err = multiply(12)
fmt.Println(<-nums, err)
// 120 <nil>
err = multiply(13)
err = multiply(14)
fmt.Println(err)
// busy
}
110 <nil>
120 <nil>
busy
Pipelines
Pipeline — это последовательность операций, где каждый шаг принимает входные данные, обрабатывает их определённым образом и выдаёт результат. Входом и выходом каждой операции является канал.
Типичный pipeline выглядит так:
- Reader: Читает входные данные из файла, базы данных или сети.
- N процессоров: Трансформируют, фильтруют, агрегируют или обогащают данные с использованием внешних источников.
- Writer: Записывает обработанные данные в файл, базу данных или сеть.
func read[T any]() <-chan T {
out := make(chan T)
go func() {
defer close(out)
for {
// читаем данные откуда-то
data := // ...
out <- data
}
}()
return out
}
func process[T any](in <-chan T) <-chan T {
out := make(chan T)
go func() {
defer close(out)
for inData := range in {
// обрабатываем данные
outData = // ...
out <- outData
}
}()
return out
}
func write[T any](in <-chan T) <-chan struct{} {
done := make(chan struct{})
go func() {
defer close(done)
for data := range in {
// записываем данные
}
}()
return done
}
Выходной канал
Горутина может сигнализировать другим горутинам о завершении работы с помощью выходного канала:
func generate(start, stop int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for i := start; i < stop; i++ {
out <- i
}
}()
return out
}
func main() {
nums := generate(5, 10)
for n := range nums {
fmt.Print(n, " ")
}
}
5 6 7 8 9
Done канал
Если горутина не должна возвращать результаты, может сигнализировать о завершении с помощью done канала:
func work() <-chan struct{} {
done := make(chan struct{})
go func() {
defer close(done)
fmt.Println("work done")
}()
return done
}
func main() {
done := work()
<-done
}
work done
Cancel канал
Для досрочного прекращения работы горутины вызывающая горутина может использовать cancel канал:
func generate(cancel chan struct{}, n int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for i := 1; i <= n; i++ {
select {
case out <- i:
case <-cancel:
return
}
}
}()
return out
}
func main() {
cancel := make(chan struct{})
defer close(cancel)
nums := generate(cancel, 10)
fmt.Println(<-nums)
fmt.Println(<-nums)
fmt.Println(<-nums)
}
1
2
3
Обработка ошибок
В параллельных pipelines есть три подхода к обработке ошибок.
➊ Возврат при первой ошибке:
// calculate выдаёт ответы для заданных чисел.
func process(in <-chan int) (<-chan int, <-chan error) {
out := make(chan Answer)
errc := make(chan error, 1)
go func() {
defer close(out)
for n := range in {
ans, err := fetchAnswer(n)
if err != nil {
errc <- err // возврат с ошибкой
return
}
out <- ans
}
errc <- nil // возврат без ошибки
}()
return out, errc
}
➋ Использование типа результата:
// Result содержит ответ или ошибку.
type Result struct {
answer int
err error
}
// calculate выдаёт ответы для заданных чисел.
func calculate(in <-chan int) <-chan Result {
out := make(chan Result)
go func() {
defer close(out)
for n := range in {
ans, err := fetchAnswer(n)
out <- Result{ans, err} // возврат ответа + ошибка
}
}()
return out
}
➌ Сбор ошибок отдельно:
// calculate выдаёт ответы для заданных чисел.
func calculate(in <-chan int, errc chan<- error) <-chan int {
out := make(chan Answer)
go func() {
defer close(out)
for n := range in {
ans, err := fetchAnswer(n)
if err == nil {
out <- ans // отправка ответа
} else {
errc <- err // или ошибка
}
}
}()
return out
}
Время
Помимо работы с датами и временем, пакет time предоставляет инструменты для управления операциями, чувствительными ко времени, в параллельных программах.
After
time.After() возвращает канал, который сначала пуст, но получает значение через период ожидания. Полезно для таймаутов операций:
// withTimeout выполняет функцию с заданным таймаутом.
func withTimeout(timeout time.Duration, fn func()) error {
done := make(chan struct{})
go func() {
defer close(done)
fn()
}()
// блокируется до завершения fn или истечения таймаута,
// в зависимости от того, что произойдёт раньше
select {
case <-done:
return nil
case <-time.After(timeout):
return errors.New("timeout")
}
}
withTimeout() ждёт завершения fn(), но благодаря time.After() не будет ждать дольше, чем timeout:
func main() {
var err error
// завершается вовремя
err = withTimeout(
50*time.Millisecond,
func() { fmt.Println("work done") },
)
fmt.Println("err =", err)
// отменяется при таймауте
err = withTimeout(
50*time.Millisecond,
func() {
time.Sleep(100 * time.Millisecond)
fmt.Println("work done")
},
)
fmt.Println("err =", err)
}
work done
err = <nil>
err = timeout
Timer
Таймер (time.Timer) — это структура с каналом C, в который отправляется текущее время при срабатывании (истечении). Таймеры полезны для планирования будущих выполнений:
done := make(chan struct{})
timer := time.NewTimer(50 * time.Millisecond)
go func() {
eventTime := <-timer.C // блокируется на 50ms
fmt.Println("work done at", eventTime)
close(done)
}()
<-done
work done at 2009-11-10 23:00:00.05
Stop() останавливает таймер и возвращает true, если он ещё не истёк, и false в противном случае:
// таймер истекает через 50ms
timer := time.NewTimer(50 * time.Millisecond)
go func() {
eventTime := <-timer.C
fmt.Println("work done at", eventTime)
}()
// через 10ms таймер ещё не истёк
time.Sleep(10 * time.Millisecond)
if timer.Stop() {
fmt.Println("execution canceled")
} else {
fmt.Println("too late to cancel")
}
execution canceled
Часто удобнее использовать функцию-обёртку time.AfterFunc(). Она ждёт длительность d, а затем выполняет функцию f:
done := make(chan struct{})
work := func() {
fmt.Println("work done")
close(done)
}
// выполняет work через 50ms
time.AfterFunc(50*time.Millisecond, work)
<-done
work done
time.AfterFunc() возвращает таймер, который можно отменить до начала выполнения:
// выполняет функцию через 50ms
timer := time.AfterFunc(50*time.Millisecond, func() {})
// через 10ms таймер ещё не истёк
time.Sleep(10 * time.Millisecond)
if timer.Stop() {
fmt.Println("execution canceled")
}
execution canceled
Если таймер используется в цикле, лучше создать один таймер и сбросить его на каждой итерации, вместо создания нового экземпляра:
// consumer читает токены из входного канала и выдаёт предупреждение,
// если значение не появляется в канале в течение часа.
func consumer(in <-chan token) {
const timeout = time.Hour
timer := time.NewTimer(timeout)
for {
timer.Reset(timeout)
select {
case <-in:
// делаем что-то
case <-timer.C:
// логируем предупреждение
}
}
}
// Допустим, отправляем 10,000 значений в канал in
// и измеряем использование памяти.
Memory used: 4 KB, # allocations: 6
Ticker
Ticker подобен таймеру, но срабатывает повторно до остановки. Ticker'ы полезны для выполнения периодических задач:
// срабатывает каждые 50ms
ticker := time.NewTicker(50 * time.Millisecond)
defer ticker.Stop()
go func() {
for {
// ждёт срабатывания ticker на каждой итерации
at := <-ticker.C
fmt.Println("work done at", at)
}
}()
// достаточно времени для срабатывания ticker 3 раза
time.Sleep(160*time.Millisecond)
ticker.Stop()
work done at 2009-11-10 23:00:00.05
work done at 2009-11-10 23:00:00.10
work done at 2009-11-10 23:00:00.15
NewTicker(d) создаёт ticker, который отправляет текущее время в канал C с интервалом d. Необходимо в конце остановить ticker с помощью Stop(), чтобы освободить ресурсы.
Если читатель канала не справляется с ticker'ом, тот будет пропускать срабатывания.
Context
Основная цель context — отмена операций, либо вручную, либо по таймауту/дедлайну.
Функция принимает context и использует его канал Done() для прослушивания отмены:
// work выполняет задачу на протяжении 50 ms, если не отменена.
// Возвращает ошибку при отмене.
func work(ctx context.Context) error {
done := make(chan struct{})
go func() {
time.Sleep(50 * time.Millisecond)
fmt.Println("work done")
close(done)
}()
select {
case <-done:
return nil
case <-ctx.Done():
return ctx.Err()
}
}
Отмена вручную (context.Canceled ошибка):
func main() {
// пустой context
ctx := context.Background()
// контекст с ручной отменой
ctx, cancel := context.WithCancel(ctx)
// отмена работает
go work(ctx)
time.Sleep(25 * time.Millisecond)
cancel()
time.Sleep(100 * time.Millisecond)
}
Отмена по таймауту (context.DeadlineExceeded ошибка):
func main() {
ctx := context.Background()
// контекст с таймаутом
ctx, cancel := context.WithTimeout(ctx, 50*time.Millisecond)
defer cancel()
// превышает таймаут
go work(ctx)
time.Sleep(100 * time.Millisecond)
}
Отмена по дедлайну:
func main() {
ctx := context.Background()
// контекст с дедлайном
deadline := time.Now().Add(50 * time.Millisecond)
ctx, cancel := context.WithDeadline(ctx, deadline)
defer cancel()
// превышает дедлайн
go work(ctx)
time.Sleep(100 * time.Millisecond)
}
Wait groups
Мы уже видели sync.WaitGroup в примере с горутинами. Рассмотрим его более детально.
Wait group помогает дождаться завершения всех горутин в группе:
func main() {
var wg sync.WaitGroup
wg.Add(2)
go func() {
defer wg.Done()
fmt.Println("worker 1")
}()
go func() {
defer wg.Done()
fmt.Println("worker 2")
}()
wg.Wait()
}
Недостаток: если горутина паникует, Wait() не знает об этом.
Альтернатива с контролем ошибок:
func main() {
ctx, cancel := context.WithCancel(context.Background())
var wg sync.WaitGroup
// запускаем горутину и следим за ошибками
wg.Add(1)
go func() {
defer wg.Done()
if err := someTask(); err != nil {
cancel()
}
}()
wg.Wait()
}
Data races
Data race происходит, когда две горутины одновременно получают доступ к одной переменной, и хотя бы одна из них записывает в неё.
Пример race condition:
var count int
go func() {
for i := 0; i < 1000; i++ {
count++ // горутина 1 читает и записывает
}
}()
go func() {
for i := 0; i < 1000; i++ {
count++ // горутина 2 читает и записывает
}
}()
time.Sleep(100 * time.Millisecond)
fmt.Println(count) // может быть меньше 2000
Go включает -race флаг для выявления race conditions при тестировании:
go test -race ./...
Race conditions
Race condition — это ошибка в коде, вызванная race condition в доступе к данным. Race condition может привести к непредсказуемому поведению программы.
Типичный пример:
var result string
go func() {
result = "done"
}()
time.Sleep(10 * time.Millisecond)
fmt.Println(result) // не гарантированно "done"
Mutexes
Mutex (mutual exclusion) защищает переменную от одновременного доступа. Только одна горутина может удерживать блокировку в один момент времени:
var (
count int
mu sync.Mutex
)
go func() {
for i := 0; i < 1000; i++ {
mu.Lock()
count++
mu.Unlock()
}
}()
go func() {
for i := 0; i < 1000; i++ {
mu.Lock()
count++
mu.Unlock()
}
}()
time.Sleep(100 * time.Millisecond)
fmt.Println(count) // всегда 2000
Lock() блокирует mutex. Если он уже заблокирован, горутина ждёт. Unlock() разблокирует его, позволяя другой горутине захватить блокировку.
Хорошая практика — использовать defer для гарантированного разблокирования:
mu.Lock()
defer mu.Unlock()
// критическая секция
RWMutex — читай-писательский mutex. Несколько горутин могут одновременно читать, но писать может только одна:
var (
data string
mu sync.RWMutex
)
// чтение
mu.RLock()
defer mu.RUnlock()
_ = data
// запись
mu.Lock()
defer mu.Unlock()
data = "new value"
Semaphores
Семафор — это счётчик, который позволяет N горутинам одновременно получить доступ к ресурсу. В Go семафоры реализуются с помощью буферизованных каналов:
// семафор с ёмкостью 3
sem := make(chan struct{}, 3)
// захватываем слот
sem <- struct{}{}
// освобождаем слот
<-sem
Типичное использование для ограничения параллелизма:
sem := make(chan struct{}, 3) // максимум 3 одновременно
var wg sync.WaitGroup
for i := 0; i < 10; i++ {
wg.Add(1)
go func() {
defer wg.Done()
sem <- struct{}{} // захватываем
defer func() { <-sem }() // освобождаем
// выполняем работу
}()
}
wg.Wait()
Signaling
Signaling — способ отправить сигнал из одной горутины в другую через канал.
Простой сигнал (без данных):
signal := make(chan struct{})
go func() {
time.Sleep(50 * time.Millisecond)
signal <- struct{}{} // посылаем сигнал
}()
<-signal // ждём сигнала
fmt.Println("done")
Сигнал с данными:
signal := make(chan string)
go func() {
signal <- "hello"
}()
msg := <-signal
fmt.Println(msg)
Run once
sync.Once гарантирует, что функция выполнится только один раз, даже если вызвана из множества горутин:
var once sync.Once
var instance string
once.Do(func() {
instance = "initialized"
})
// остальные вызовы будут проигнорированы
once.Do(func() {
instance = "should not run"
})
fmt.Println(instance) // "initialized"
Полезно для ленивой инициализации:
var (
instance *Database
once sync.Once
)
func GetDatabase() *Database {
once.Do(func() {
instance = &Database{}
instance.Connect()
})
return instance
}
Object pool
Object pool — паттерн для переиспользования объектов и снижения нагрузки на сборщик мусора:
pool := sync.Pool{
New: func() interface{} {
return &bytes.Buffer{}
},
}
// получаем объект из пула
buf := pool.Get().(*bytes.Buffer)
defer pool.Put(buf)
// используем объект
buf.WriteString("hello")
Atomics
Atomic операции позволяют безопасно получать доступ к простым значениям без mutex:
var count atomic.Int64
// атомично увеличиваем
count.Add(1)
// атомично получаем значение
val := count.Load()
// атомично записываем
count.Store(10)
Атомики быстрее mutex'ов для простых операций, но работают только с базовыми типами.
Testing
Тестирование параллельного кода сложнее, так как поведение недетерминировано.
Используйте флаг -race:
go test -race ./...
Для тестирования таймаутов используйте time.Sleep или context:
func TestTimeout(t *testing.T) {
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
defer cancel()
if err := work(ctx); err != context.DeadlineExceeded {
t.Fatalf("expected timeout, got %v", err)
}
}
Для повторяемости тестов используйте синхронизацию через каналы вместо time.Sleep.
Scheduling
runtime.Gosched() позволяет текущей горутине добровольно сдать управление другим горутинам:
import "runtime"
for i := 0; i < 10; i++ {
fmt.Println(i)
runtime.Gosched() // даём возможность другим горутинам работать
}
В большинстве случаев это не требуется, так как scheduler работает автоматически.
Diagnostics
Go предоставляет инструменты для анализа параллельного кода.
runtime.NumGoroutine() возвращает количество горутин:
fmt.Println(runtime.NumGoroutine()) // количество активных горутин
Профилирование dengan pprof:
import _ "net/http/pprof"
go func() {
log.Println(http.ListenAndServe("localhost:6060", nil))
}()
// посетите http://localhost:6060/debug/pprof/
Трассировка с помощью go trace:
f, _ := os.Create("trace.out")
defer f.Close()
trace.Start(f)
defer trace.Stop()
// ваш код
Заключение
Параллелизм в Go — мощный инструмент для написания масштабируемых приложений. Основные концепции:
- Горутины — лёгкие потоки управления, легко создавать сотни.
- Каналы — безопасный способ обмена данными между горутинами.
- Select — управление множественными каналами.
- Context — управление отменой и таймаутами.
- Синхронизация — mutex, wait groups, atomic операции для защиты данных.
При разработке параллельного кода всегда проверяйте наличие race conditions с флагом -race и тщательно тестируйте.