Поток обработки и управление его временем жизни
Pipeline — цепочка этапов: один производит значения, следующие читают их, преобразуют и отправляют дальше. Канал связывает выход предыдущего этапа со входом следующего. Fan-in соединяет несколько входных потоков в один выход.
У каждой стадии должен быть понятный владелец её выходного канала. Производитель закрывает собственный выход после последней отправки. Потребитель не закрывает канал, в который сам не отправляет: он не может знать, завершили ли работу остальные отправители. Для fan-in выход закрывается только после того, как все копирующие горутины закончили пересылку.
Если потребитель закончил раньше источников, заблокированные отправители могут никогда не проснуться. Передавай context.Context в каждую стадию и выбирай между отправкой и ctx.Done(). Отмену инициирует вызывающая сторона; после неё стадии перестают ждать друг друга, закрывают свои выходы и завершаются. Если поток полностью прочитан, каналы закрываются естественно. Не используй time.Sleep как замену ожиданию или отмене.
Пример
package main
import (
"context"
"fmt"
"sync"
)
func generate(ctx context.Context, values ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, value := range values {
select {
case out <- value:
case <-ctx.Done():
return
}
}
}()
return out
}
func square(ctx context.Context, input <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for {
select {
case value, ok := <-input:
if !ok {
return
}
select {
case out <- value * value:
case <-ctx.Done():
return
}
case <-ctx.Done():
return
}
}
}()
return out
}
func fanIn(ctx context.Context, inputs ...<-chan int) <-chan int {
out := make(chan int)
var workers sync.WaitGroup
workers.Add(len(inputs))
for _, input := range inputs {
go func(input <-chan int) {
defer workers.Done()
for {
select {
case value, ok := <-input:
if !ok {
return
}
select {
case out <- value:
case <-ctx.Done():
return
}
case <-ctx.Done():
return
}
}
}(input)
}
go func() {
workers.Wait()
close(out)
}()
return out
}
func main() {
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
leftInput := generate(ctx, 1, 2, 3, 4)
rightInput := generate(ctx, 1, 2, 3, 4)
left := square(ctx, leftInput)
right := square(ctx, rightInput)
merged := fanIn(ctx, left, right)
first := <-merged
cancel()
for range merged {
// Drain until all fan-in workers have exited and closed the output.
}
for range left {
}
for range right {
}
for range leftInput {
}
for range rightInput {
}
fmt.Printf("первый результат: %d; pipeline остановлен\n", first)
}
Оба источника отправляют одинаковую последовательность, поэтому какой бы поток ни пришёл первым, первое значение после возведения в квадрат равно 1. Потребитель отменяет контекст после этого результата и продолжает читать выход fan-in до его закрытия. Это ожидание важно: оно подтверждает, что пересыльщики завершились.
Каждый генератор и преобразователь закрывает только собственный выход. Fan-in ждёт всех пересыльщиков через группу ожидания, затем закрывает общий канал. Отмена одновременно может разблокировать производителей на отправке и потребителей на получении. Пример специально не печатает следующие результаты: их количество после конкурентной отмены зависит от того, какие операции успели начаться.
Для полного прохода по данным потребитель обычно читает канал до закрытия. Если же он решает остановиться раньше, он должен отменить контекст, иначе верхние стадии могут остаться заблокированными при отправке.