← Все этапы
05 · КОНКУРЕНТНОСТЬУрок 33 из 3328 минут

Pipeline, fan-in и отмена

Соединяем этапы каналами, сводим потоки в один и останавливаем работу через context.

Поток обработки и управление его временем жизни

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 ждёт всех пересыльщиков через группу ожидания, затем закрывает общий канал. Отмена одновременно может разблокировать производителей на отправке и потребителей на получении. Пример специально не печатает следующие результаты: их количество после конкурентной отмены зависит от того, какие операции успели начаться.

Для полного прохода по данным потребитель обычно читает канал до закрытия. Если же он решает остановиться раньше, он должен отменить контекст, иначе верхние стадии могут остаться заблокированными при отправке.

Практика

Попробуйте сами

Построй два генератора, один с числами 2, 4, 6, второй с 1, 3, 5. Пропусти каждый поток через стадию, которая умножает число на 10, затем объедини их через fan-in. Используй контекст отмены и корректное закрытие выходов. В основной функции полностью прочитай поток, собери значения, отсортируй и напечатай.

Проверка результата

Как понять, что получилось

Вывод всегда равен [10 20 30 40 50 60]. Каждый генератор закрывает свой канал; каждая преобразующая стадия закрывает свой выход; fan-in закрывает общий выход только после завершения всех пересыльщиков. Все отправки и ожидания стадии могут завершиться по ctx.Done(), а программа не использует задержки для синхронизации.

Pipeline, fan-in и отмена | Go | WebSchool · Go