СОБЕСЕДОВАНИЕ в КАСПЕРСКИЙ. GOLANG РАЗРАБОТКА
Сегодня мы разберём, как кандидат на позицию 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 / Жесткая остановка)
- Меняем состояние на
Stoppingпод мьютексом. - Дренируем очередь: извлекаем все задачи, отвечаем им ошибкой
ErrPoolStopping(дляSubmitWaitпишем в ихdoneканал). - Будим всех воркеров (
Broadcast), чтобы они увидели новое состояние и пустую очередь. - Ждем
wg.Wait()— завершения только активных задач. - Переводим в
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 канала нельзя отличить "нет задач" от "пул остановлен", Submit будет паниковать при отправке в закрытый канал. Нужно дополнительное состояние (atomic/mutex) для проверки перед отправкой.", "correct_answer": ""}
