Реклама
Перетяжка // Коробка 3.0

Producer-Consumer на Java, Go и Rust: одна задача и три уровня решения

Задачу с собеседования почти все решают ручной координацией потоков, хотя у неё есть три разных уровня решения. Разбираем, чем они отличаются и по какому признаку выбирается нужный.

Обложка: Producer-Consumer на Java, Go и Rust: одна задача и три уровня решения

Задача Producer-Consumer звучит обманчиво просто: одни потоки производят элементы, другие их разбирают, а между ними стоит буфер. На собеседовании её просят решить почти всегда, и почти всегда кандидат берётся писать координацию потоков руками.

Между тем у этой задачи есть три принципиально разных уровня решения, и выбор уровня определяется не языком, а требованиями к задержке. Разбираем все три: блокирующую очередь на Java, каналы и горутины в Go и очередь без блокировок на Rust.

Producer-Consumer — это схема, в которой производители и потребители работают независимо, обмениваясь данными через общий буфер конечной ёмкости. Ёмкость здесь работает главным механизмом защиты: именно она не даёт быстрым производителям завалить медленных потребителей.

Ключевые выводы

Ручная координация через wait() и notifyAll() в блоках синхронизации даёт гонки и ложные пробуждения, а активное ожидание в цикле впустую жжёт процессор.

Готовая блокирующая очередь решает задачу целиком: при переполнении она сама останавливает производителя, обеспечивая обратное давление вместо переполнения памяти.

Горутины — это не потоки операционной системы, а лёгкие единицы, которые планировщик рантайма раскладывает по потокам, поэтому создавать их можно в количествах, немыслимых для системных потоков.

Три гарантии прогресса различаются строго: obstruction-free требует отсутствия помех, lock-free гарантирует прогресс системы в целом, wait-free — завершение каждой операции за ограниченное число шагов.

Отсутствие блокировок и использование атомарных операций сами по себе не дают wait-free: цикл повторных попыток без верхней границы — это lock-free, как бы его ни называли в описании библиотеки.

Уровень первый: блокирующая очередь Producer-Consumer

С этого уровня начинают почти все, и здесь же делают три типовые ошибки:

  • Писать собственную логику блокировок через wait() и notifyAll() внутри синхронизированных блоков. Это приносит тонкие гонки и баги ложного пробуждения, когда поток просыпается без всякого уведомления.
  • Крутить активное ожидание в бесконечном цикле поверх коллекции, не рассчитанной на многопоточность. Процессор загружен постоянной проверкой размера очереди, а корректности всё равно нет.
  • Не предусмотреть обратного давления. Когда производители быстрее потребителей, буфер растёт, пока приложение не упрётся в нехватку памяти.

Правильная ментальная модель здесь кухонная: повара выкладывают блюда на раздачу фиксированной вместимости, официанты забирают их по готовности. Раздача заполнена — повар ждёт, и это нормальная работа системы, а не сбой.

В Java всё это уже реализовано в стандартной библиотеке:

			public class KitchenPass {
    private final BlockingQueue<Order> pass = new ArrayBlockingQueue<>(10);

    public void prepareOrder(Order order) throws InterruptedException {
        pass.put(order);      // блокируется, если раздача заполнена
    }

    public Order collectOrder() throws InterruptedException {
        return pass.take();   // блокируется, если раздача пуста
    }
}
		

Вся координация потоков спрятана внутри реализации очереди, и явных блокировок в коде не остаётся вовсе. Метод put останавливает производителя при переполнении, и это и есть обратное давление, ради которого затевалась ограниченная ёмкость.

Уровень второй: передача сообщений вместо общей памяти

Go подходит к задаче с другой стороны. Вместо того чтобы защищать общий буфер, язык предлагает передавать данные между независимыми исполнителями. Запускается такой исполнитель одним словом:

			func sayHello() {
    fmt.Println("Привет из горутины")
}

func main() {
    go sayHello()
    fmt.Println("Привет из main")
}
		

Ключевое слово go перед вызовом запускает функцию конкурентно. Правда, у примера выше есть подвох: когда main завершается, программа выходит целиком, и горутина может не успеть отработать. Дожидаться её нужно явно:

			var wg sync.WaitGroup
wg.Add(2)

go func() {
    defer wg.Done()
    fmt.Println("Задача 1 завершена")
}()
go func() {
    defer wg.Done()
    fmt.Println("Задача 2 завершена")
}()

wg.Wait()
fmt.Println("Все задачи завершены")
		

Где здесь, собственно, очередь

Горутины сами по себе задачу производителя и потребителя не решают: WaitGroup лишь дожидается завершения, буфера и обратного давления в нём нет. Роль ограниченной очереди в Go играет буферизированный канал, и это прямой аналог кухонной раздачи из предыдущего раздела:

			var pass = make(chan Order, 10)   // та же вместимость, что у ArrayBlockingQueue

func prepareOrder(order Order) {
    pass <- order              // блокируется, если канал заполнен
}

func collectOrder() Order {
    return <-pass              // блокируется, если канал пуст
}
		

Поведение совпадает с блокирующей очередью один в один: отправка встаёт при переполнении, приём встаёт на пустом канале. Разница в акценте. В Java вы защищаете общий буфер, в Go передаёте значение из одной горутины в другую, и владение данными переходит вместе с ним.

Почему это не просто потоки с коротким именем

Горутина не является потоком операционной системы. Это лёгкая единица конкурентного исполнения, которой управляет рантайм языка: он раскладывает горутины по небольшому набору системных потоков сам. Отсюда и практическое следствие — можно оперировать очень большим числом конкурентных задач, не создавая и не сопровождая отдельный системный поток под каждую.

Здесь же полезно развести два понятия, которые постоянно путают. Конкурентность — это способ структурировать программу так, чтобы несколько задач продвигались независимо. Параллелизм означает фактическое одновременное исполнение на нескольких вычислительных ресурсах. Программа может быть конкурентной, и при этом ни одна пара её задач не выполняется буквально в один момент.

Уровень третий: очередь без блокировок

Третий уровень нужен, когда неприемлема сама возможность того, что поток уснёт на блокировке. Это низкая задержка в торговых системах, обработка звука, ядра планировщиков. И здесь начинается терминология, в которой путаются чаще всего.

Три гарантии прогресса

Obstruction-free. Операция завершится, если в какой-то момент ей дадут поработать без помех со стороны других потоков. При непрерывных помехах прогресс не гарантирован вообще. Самая слабая из трёх гарантий.

Lock-free. Гарантирует прогресс системы в целом: даже если один поток застрял, какой-то другой обязательно завершит свою операцию. Отдельный поток при этом может голодать сколь угодно долго, но система не встаёт целиком.

Wait-free. Каждая операция завершается за ограниченное число шагов независимо от того, что делают остальные потоки. Ключевая формулировка здесь именно «ограниченное число шагов», а не «обычно быстро», «не использует блокировок» или «работает на атомарных операциях». Это формальное свойство, а не характеристика скорости.

Два неравенства, которые стоит запомнить:
Отсутствие блокировок не равно wait-free. Использование атомарных операций тоже не равно wait-free.

Почему мьютекса недостаточно

Очередь, обёрнутая в мьютекс, — это совершенно корректный многопоточный код на Rust. Но wait-free она не является: поток, захвативший блокировку, может быть снят с процессора планировщиком, и все остальные будут ждать его возвращения неограниченно долго. Для большинства задач это приемлемо, для систем с жёсткими требованиями по задержке уже нет.

Атомарность и упорядочивание — разные вещи

Процессор гарантирует, что атомарная операция неделима. Этого мало: нужно ещё, чтобы потребитель увидел записи, сделанные производителем до публикации флага. За это отвечает модель упорядочивания памяти.

			// производитель
buffer[index] = value;
ready.store(true, Ordering::Release);

// потребитель
if ready.load(Ordering::Acquire) {
    let value = buffer[index];   // видит запись, сделанную до Release
}
		

Пара «освобождение и захват» создаёт отношение синхронизации: если потребитель увидел освобождающую запись, он гарантированно видит и все предшествовавшие ей записи. Это одна из центральных идей программирования без блокировок, и именно её упускают, заменяя мьютекс на атомарную операцию без всего остального.

Кольцевой буфер и номера последовательности

Практическая конструкция ограниченной очереди для нескольких производителей и нескольких потребителей строится на кольцевом буфере с двумя атомарными счётчиками логических позиций. Физический слот вычисляется как остаток от деления позиции на ёмкость, а при ёмкости, равной степени двойки, деление заменяется маской по битам.

Каждый слот хранит не только значение, но и номер последовательности, по которому определяется, какому логическому поколению этот слот сейчас принадлежит:

			struct Slot<T> {
    sequence: AtomicUsize,
    value: UnsafeCell<MaybeUninit<T>>,
}
		

Тип UnsafeCell здесь означает ровно то, что написано на упаковке: обращаться к значению можно только из блока unsafe, компилятор проверку прекращает, и инвариант «в конкретный момент со слотом работает один поток» держится на номере последовательности, а не на системе типов.

Тип MaybeUninit здесь не для красоты: при ёмкости в 1024 элемента нужны 1024 места под значения, а не 1024 сконструированных объекта. Производитель получает уникальную позицию одним атомарным инкрементом, и никакого соревнования за неё не возникает:

			let position = self.enqueue_pos.fetch_add(1, Ordering::Relaxed);
let index = position & self.mask;

let slot = &self.buffer[index];
while slot.sequence.load(Ordering::Acquire) != position {
    std::hint::spin_loop();      // цикл ожидания без верхней границы
}
		

Последние три строки стоит запомнить, потому что через абзац они окажутся важны. Позиция достаётся производителю без соревнования, а вот дальше он крутится в ожидании, пока слот не перейдёт к его поколению. Верхней границы у этого ожидания нет.

Где заканчивается честность в названиях

Неприятная правда, которую стоит знать до выбора библиотеки: многое из того, что продаётся как wait-free очередь, на деле является lock-free. Причина в устройстве ядра алгоритма:

			loop {
    let pos = tail.load(...);
    if try_reserve(pos) { break; }
}
		

Операция обычно завершается быстро, но верхней границы на число неудачных попыток здесь нет, а именно она и определяет wait-free. И да, ровно такой цикл вы только что видели в нашей собственной очереди: строго говоря, на этом шаге она тоже lock-free, а не wait-free, и автор исходного разбора это прямо оговаривает. Строгая реализация требует ограниченной стратегии резервирования или механизма помощи: когда поток B встречает незавершённую операцию потока A, он не ждёт возвращения A, а доводит его операцию до конца сам. Так прогресс перестаёт зависеть от того, что планировщик сделал с конкретным потоком.

Какой уровень выбрать

Правило простое и почти всегда работает от обратного.

  • Блокирующая очередь из стандартной библиотеки закрывает подавляющее большинство задач. Начинайте с неё и не пишите координацию потоков вручную.
  • Модель с передачей сообщений выигрывает, когда конкурентных задач много, а связи между ними хочется описать явно, а не через набор общих структур под блокировками.
  • Алгоритмы без блокировок берутся, когда есть измеренное требование к задержке, которое блокировка нарушает: не «хочется быстрее», а конкретная цифра в бюджете. И в продакшене их берут готовыми из библиотек вроде crossbeam, а разбирают по косточкам ради понимания механики.

Ошибка почти всегда идёт в одну сторону: за атомарные операции берутся раньше, чем появилась причина, и получают код, который сложнее отлаживать, а гарантий даёт не больше.

Часто задаваемые вопросы
1
Что такое паттерн Producer-Consumer?

Схема, в которой одни потоки создают элементы работы, другие их обрабатывают, а между ними стоит общий буфер конечной ёмкости. Ограниченная ёмкость обеспечивает обратное давление: когда буфер заполнен, производитель приостанавливается вместо того, чтобы копить данные в памяти.

2
Почему не стоит писать координацию потоков через wait и notifyAll?

Ручная реализация даёт тонкие гонки и уязвима к ложным пробуждениям, когда поток просыпается без уведомления. Готовая блокирующая очередь инкапсулирует всю координацию, и явных блокировок в прикладном коде не остаётся.

3
Чем горутина отличается от потока операционной системы?

Горутина — это лёгкая единица исполнения, которой управляет рантайм языка: он сам распределяет горутины по небольшому набору системных потоков. Поэтому создать десятки тысяч горутин можно, а десятки тысяч системных потоков нет.

4
В чём разница между lock-free и wait-free?

Lock-free гарантирует прогресс системы в целом: кто-то из потоков обязательно завершит операцию, но конкретный поток может голодать. Wait-free сильнее: каждая операция завершается за ограниченное число шагов независимо от поведения остальных потоков.

5
Достаточно ли атомарных операций, чтобы получить wait-free алгоритм?

Нет. Цикл повторных попыток на атомарных операциях без верхней границы числа неудач остаётся lock-free. Для строгой wait-free нужна ограниченная стратегия резервирования либо механизм помощи, при котором один поток доводит до конца незавершённую операцию другого.

6
Зачем в очереди без блокировок нужен порядок памяти Release и Acquire?

Атомарность гарантирует неделимость самой операции, но не гарантирует, что потребитель увидит записи, сделанные производителем до публикации флага. Пара освобождения и захвата создаёт отношение синхронизации, при котором предшествующие записи становятся видимы.

Что забрать с собой

Задача Producer-Consumer решается на трёх уровнях, и уровень выбирается требованиями, а не вкусом. Ограниченный буфер с обратным давлением закрывает большинство случаев. Передача сообщений упрощает структуру, когда конкурентных задач становится много. Алгоритмы без блокировок нужны там, где блокировка нарушает измеренный бюджет задержки, и тогда важно понимать, какую именно гарантию вы покупаете.

Материалы разбора: задача о ресторанной кухне на Java, введение в горутины и модель конкурентности Go и построение wait-free очереди на Rust.

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