Перейти к основному содержимому

СОБЕСЕДОВАНИЕ в КАСПЕРСКИЙ. GOLANG РАЗРАБОТКА

· 8 мин. чтения

Сегодня мы разберём, как кандидат на позицию Go‑разработчика в Касперском решает задачу реализации worker‑pool с неблокирующим submit, мягкой и жёсткой остановкой, а также как он последовательно переходит от базового решения на слайсе и мьютексе к более изящному варианту, использующему sync.Cond. В ходе разбора показаны ключевые детали конкурентного программирования: управление состоянием пула, взаимодействие воркеров через каналы и правильное завершение горутин. Это собеседование демонстрирует не только знание синхронных примитивов Go, но и способность кандидата articulate trade‑offs и улучшать код в реальном времени.

Вопрос 1. Реализовать интерфейс worker pool с методами: NewWorkerPool (конструктор), Submit (неблокирующее добавление задачи), SubmitWait (добавление задачи с ожиданием выполнения), Stop (жёсткая остановка: дождаться только выполняющихся задач), StopWait (мягкая остановка: дождаться всех задач включая очередь). Обеспечить соблюдение порядка очереди, неблокирующий Submit, корректную обработку ошибок и состояний.

Таймкод: 00:00:23

Ответ собеседника: Правильный. Предложено решение на основе слайса с мьютексом для очереди задач, канала для пробуждения воркеров (WakeCh), WaitGroup для отслеживания завершения, машины состояний (Running, Stopping, Finishing, Stopped). Submit добавляет задачу под мьютексом и сигнализирует воркеру через неблокирующую отправку в WakeCh. SubmitWait создаёт буферизированный канал done, ждёт результат. Stop переводит в Stopping, забирает очередь, возвращает ошибку Stop задачам из очереди, ждёт активных. StopWait переводит в Finishing, дожидается всей очереди.

Правильный ответ:

Интерфейс и типы ошибок

Начнём с определения контракта. Использование интерфейса для задач позволяет пулу быть универсальным. Введём сентинельные ошибки для отличия причин отмены.

type Task func() error

type WorkerPool interface {
Submit(task Task) error
SubmitWait(task Task) error
Stop() error
StopWait() error
}

var (
ErrPoolStopped = errors.New("worker pool: pool is stopped")
ErrPoolStopping = errors.New("worker pool: pool is stopping, queued tasks rejected")
ErrTaskRejected = errors.New("worker pool: task rejected")
)

Состояния пула (State Machine)

Четкое разделение состояний критично для корректной работы Stop и StopWait без гонок (race conditions).

  • Running — нормальная работа. Принимаем задачи в очередь, воркеры обрабатывают.
  • Stopping — вызван Stop(). Новые задачи в Submit/SubmitWait отклоняются (ErrPoolStopping). Задачи из очереди извлекаются и завершаются с ошибкой ErrPoolStopping. Активные воркеры дорабатывают.
  • Finishing — вызван StopWait(). Новые задачи отклоняются (ErrPoolStopping). Очередь не дренируется, воркеры дорабатывают все задачи из очереди.
  • Stopped — все воркеры завершены, ресурсы освобождены. Любые вызовы методов возвращают ErrPoolStopped.

Структура пула

type pool struct {
mu sync.Mutex
cond *sync.Cond // для эффективного ожидания задач/состояний
tasks []taskItem // очередь (ring buffer или слайс с head/tail)
workers int // количество запущенных воркеров
state state // текущее состояние
wg sync.WaitGroup // ожидание завершения активных задач
closeOnce sync.Once // гарантия однократного закрытия ресурсов
}

type taskItem struct {
fn Task
done chan error // nil для Submit, буферизированный для SubmitWait
}

type state uint32

const (
stateRunning state = iota
stateStopping
stateFinishing
stateStopped
)

Конструктор NewWorkerPool

Запускаем фиксированное количество воркеров. Используем sync.Cond вместо отдельного WakeCh — это идиоматичнее для работы с очередью под мьютексом и позволяет избежать "пробуждения" воркеров, когда задач нет (spurious wakeups обрабатываются циклом for).

func NewWorkerPool(workers int) WorkerPool {
if workers <= 0 {
workers = runtime.NumCPU()
}
p := &pool{
tasks: make([]taskItem, 0, 64), // начальная емкость
}
p.cond = sync.NewCond(&p.mu)

p.workers = workers
p.state = stateRunning

for i := 0; i < workers; i++ {
p.wg.Add(1)
go p.worker()
}
return p
}

Основной цикл воркера (worker)

Воркер ждет задачу под мьютексом. Ключевой момент: проверка состояния после пробуждения и перед выполнением задачи.

func (p *pool) worker() {
defer p.wg.Done()

for {
p.mu.Lock()
// Ждем задачу ИЛИ изменение состояния (Stop/StopWait)
for len(p.tasks) == 0 && p.state == stateRunning {
p.cond.Wait()
}

// Если пул останавливается и очередь пуста — выходим
if (p.state == stateStopping || p.state == stateFinishing) && len(p.tasks) == 0 {
p.mu.Unlock()
return
}

// Если Running, но задач нет (спуриусный wakeup) — продолжаем цикл
if len(p.tasks) == 0 {
p.mu.Unlock()
continue
}

// Извлекаем задачу (FIFO)
item := p.tasks[0]
p.tasks = p.tasks[1:]
p.mu.Unlock()

// Выполняем задачу БЕЗ мьютекса
err := item.fn()

// Отдаем результат, если кто-то ждет (SubmitWait)
if item.done != nil {
item.done <- err
}
}
}

Метод Submit (неблокирующий)

Должен вернуть ошибку мгновенно, если пул не в состоянии Running. Используем select с default для неблокирующего уведомления воркера через cond.Signal() (хотя Broadcast безопаснее при множестве воркеров, Signal достаточно, так как одна задача = один воркер).

func (p *pool) Submit(task Task) error {
if task == nil {
return errors.New("worker pool: nil task")
}

p.mu.Lock()
defer p.mu.Unlock()

if p.state != stateRunning {
return ErrPoolStopping // или ErrPoolStopped если уже Stopped
}

p.tasks = append(p.tasks, taskItem{fn: task, done: nil})
p.cond.Signal() // Будим одного воркера
return nil
}

Метод SubmitWait (блокирующий с результатом)

Создаем буферизированный канал done (буфер 1 критически важен, чтобы воркер не блокировался на отправке результата, если вызывающий уже ушел по таймауту/контексту, хотя здесь таймаута нет, но это best practice).

func (p *pool) SubmitWait(task Task) error {
if task == nil {
return errors.New("worker pool: nil task")
}

done := make(chan error, 1) // Буфер 1!

p.mu.Lock()
if p.state != stateRunning {
p.mu.Unlock()
return ErrPoolStopping
}

p.tasks = append(p.tasks, taskItem{fn: task, done: done})
p.cond.Signal()
p.mu.Unlock()

return <-done // Блокируемся до завершения
}

Метод Stop (Hard Stop / Жесткая остановка)

  1. Меняем состояние на Stopping под мьютексом.
  2. Дренируем очередь: извлекаем все задачи, отвечаем им ошибкой ErrPoolStopping (для SubmitWait пишем в их done канал).
  3. Будим всех воркеров (Broadcast), чтобы они увидели новое состояние и пустую очередь.
  4. Ждем wg.Wait() — завершения только активных задач.
  5. Переводим в Stopped.
func (p *pool) Stop() error {
p.mu.Lock()
if p.state != stateRunning {
p.mu.Unlock()
return ErrPoolStopped // или ErrPoolStopping если уже в процессе
}
p.state = stateStopping

// Отклоняем задачи в очереди
for _, item := range p.tasks {
if item.done != nil {
item.done <- ErrPoolStopping
}
}
p.tasks = nil // Очищаем очередь

p.cond.Broadcast() // Будим всех, чтобы они проверили state и вышли
p.mu.Unlock()

p.wg.Wait() // Ждем ТОЛЬКО выполняющихся

p.mu.Lock()
p.state = stateStopped
p.mu.Unlock()
return nil
}

Метод StopWait (Graceful Stop / Мягкая остановка)

Аналогично Stop, но не трогаем очередь. Ставим состояние Finishing. Воркеры продолжают забирать задачи до тех пор, пока очередь не опустеет. После wg.Wait() переводим в Stopped.

func (p *pool) StopWait() error {
p.mu.Lock()
if p.state != stateRunning && p.state != stateStopping {
p.mu.Unlock()
return ErrPoolStopped
}

// Если уже был Stop, переходим в Finishing, чтобы доработать остатки
p.state = stateFinishing
p.cond.Broadcast()
p.mu.Unlock()

p.wg.Wait() // Ждем ВСЕХ (активных + из очереди)

p.mu.Lock()
p.state = stateStopped
p.mu.Unlock()
return nil
}

Важные детали реализации (Senior Level)

1. Порядок очереди (FIFO) и производительность Слайс p.tasks = p.tasks[1:] приводит к сдвигу памяти O(N). Для высоконагруженных систем лучше использовать Ring Buffer (кольцевой буфер) на слайсе фиксированного размера с индексами head/tail или deque из container/list (хотя список дает аллокации). Ring Buffer дает O(1) на push/pop и сохраняет порядок.

2. sync.Cond vs Channel (WakeCh) Кандидат упомянул WakeCh. Канал подходит, но требует отдельной горутины-диспетчера или сложной логики select в воркере. sync.Cond привязан к мьютексу очереди, что дает атомарность проверки "очередь пуста" + "уснуть". С каналом легко упустить пробуждение (missed wakeup) между проверкой длины очереди и чтением из канала. Cond — стандартный примитив для Producer/Consumer на мьютексах.

3. Обработка SubmitWait при остановке Если Stop вызван пока задача ждет в SubmitWait, задача из очереди извлекается, в её done канал пишется ошибка. Вызывающий SubmitWait разблокируется и получит ErrPoolStopping. Это корректное поведение "fail-fast" для ожидающих.

4. Идемпотентность Stop / StopWait Методы безопасны для повторного вызова. Второй вызов вернет ErrPoolStopped (или ErrPoolStopping если еще не завершился первый), не паникуя и не блокируясь вечно.

5. Graceful Shutdown с context.Context (Best Practice) В реальных системах Stop/StopWait часто оборачивают в контекст с таймаутом:

func (p *pool) StopWait(ctx context.Context) error {
// ... логика StopWait ...
done := make(chan struct{})
go func() { p.wg.Wait(); close(done) }()
select {
case <-done: return nil
case <-ctx.Done(): return ctx.Err() // Таймаут graceful shutdown
}
}

6. Отсутствие утечек горутин

  • wg.Add(1) в конструкторе, wg.Done() в defer воркера.
  • Все пути выхода из воркера (break/return) покрыты defer.
  • Каналы done в SubmitWait буферизированы (или воркер не блокируется на отправке, если вызывающий уже ушел — но тут нет контекста, поэтому буфер 1 гарантирует неблокирующую отправку результата).

7. Проверка nil задачи Паника внутри воркера убьет воркер (если нет recover), уменьшив пул. Лучше валидировать на входе (Submit/SubmitWait) и/или обернуть выполнение item.fn() в func() { defer recover(); ... }().

Резюме Предложенная схема с sync.Mutex + sync.Cond + sync.WaitGroup + State Machine является классической, надежной и высокопроизводительной реализацией Worker Pool на Go. Она гарантирует FIFO, неблокирующий Submit, корректную обработку ошибок при остановке и отсутствие гонок данных{"question": "Можно ли реализовать это через channel? Если да, то как? Если нет, то почему?", "timecode": "00:14:59", "anser_status": "Правильный", "answer": "Можно. Канал — это и есть очередь. Воркеры читают из канала. Submit — неблокирующая отправка в канал (select с default). SubmitWait — блокирующая отправка или отправка структуры с каналом для ответа. Stop — close канала (жёсткая остановка: воркеры дочитают буфер и выйдут). StopWait — не закрываем канал, а ждём WaitGroup после сигнала остановки (через отдельный канал done/ctx). Проблема: при close канала нельзя отличить &#34;нет задач&#34; от &#34;пул остановлен&#34;, Submit будет паниковать при отправке в закрытый канал. Нужно дополнительное состояние (atomic/mutex) для проверки перед отправкой.", "correct_answer": ""}