Мастерство многопоточного программирования на Java: от основ до Lock-Free алгоритмов

Глубокий курс по разработке высокопроизводительных систем, охватывающий механизмы синхронизации, модель памяти JMM и современные асинхронные подходы. Вы научитесь писать эффективный и безопасный код, понимая внутреннее устройство JVM и аппаратные особенности многопроцессорных систем.

Создание потоков через Thread и Runnable: где кроются ошибки

Создание потоков через Thread и Runnable: где кроются ошибки

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

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

Что такое поток в Java?

Когда вы запускаете Java-программу, виртуальная машина (JVM) создает один главный поток — он так и называется main. Именно в нем последовательно выполняются инструкции вашего метода public static void main.

Но сам по себе класс java.lang.Thread в Java — это лишь объект-обертка. Это пульт управления. Когда вы создаете объект Thread через new, никакого реального потока в операционной системе еще не существует. Это просто кусок памяти в куче (Heap), хранящий имя потока, его приоритет и состояние.

Настоящая магия происходит только в момент запуска. JVM обращается к операционной системе и просит выделить ресурсы: системный поток, отдельный стек вызовов и регистры процессора.

Два пути создания потоков

Исторически в Java есть два способа описать код, который должен выполниться параллельно.

Путь 1: Наследование от класса Thread (Плохой подход)

Первый инстинкт — создать свой класс, унаследовать его от Thread и переопределить метод run().

public class ReportGenerator extends Thread {
    @Override
    public void run() {
        System.out.println("Генерация отчета началась в потоке: " + Thread.currentThread().getName());
        // Долгая логика генерации...
    }
}

// Использование:
ReportGenerator generator = new ReportGenerator();
generator.start();

Почему так делать не стоит?

  1. Отсутствие множественного наследования. В Java класс может наследоваться только от одного предка. Если ваш ReportGenerator уже наследует Thread, он не сможет унаследовать базовый класс BaseService из вашей бизнес-логики.
  2. Нарушение архитектуры. Класс Thread отвечает за механику управления системным потоком (приоритеты, прерывания, состояния). Ваша логика отчета — это бизнес-задача. Смешивать их в одном классе — значит нарушать принцип единственной ответственности (Single Responsibility Principle).

Путь 2: Реализация интерфейса Runnable (Правильный подход)

Чтобы отделить задачу от механизма ее выполнения, в Java существует интерфейс Runnable. В нем есть всего один метод — run().

public class ReportTask implements Runnable {
    @Override
    public void run() {
        System.out.println("Генерация отчета началась в потоке: " + Thread.currentThread().getName());
    }
}

// Использование:
ReportTask task = new ReportTask(); // Создали задачу
Thread worker = new Thread(task);   // Наняли рабочего и дали ему задачу
worker.start();                     // Сказали рабочему начать

Здесь ReportTask ничего не знает о потоках, приоритетах и операционной системе. Это просто инструкция. А класс Thread выступает в роли рабочего, которому передали эту инструкцию.

Поскольку Runnable является функциональным интерфейсом (содержит ровно один абстрактный метод), современный Java-код позволяет использовать лямбда-выражения, избавляя нас от создания отдельных классов:

Thread worker = new Thread(() -> {
    System.out.println("Анонимная задача в потоке: " + Thread.currentThread().getName());
});
worker.start();

Главная ловушка: start() против run()

Самая частая ошибка при работе с потоками выглядит так: разработчик создает объект Thread (или Runnable), описывает в нем тяжелую задачу, а затем вызывает метод run() вместо start().

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

Почему это происходит? Метод run() — это самый обычный метод Java. В нем нет никакой магии. Если вы вызываете worker.run(), вы просто говорите текущему потоку (например, main): «Выполни-ка инструкции из этого метода прямо сейчас». Текущий поток послушно идет в метод run(), выполняет его от начала до конца, и только потом переходит к следующей строчке кода.

Метод start() — это нативный вызов (Native Method). Он делает следующее:

  1. Обращается к JVM и ОС для создания нового системного потока.
  2. Подготавливает для него отдельный стек.
  3. Возвращает управление вызывающему потоку (вызывающий поток мгновенно идет дальше).
  4. А внутри нового потока JVM сама вызывает ваш метод run().

Цена создания потоков

Теперь вы умеете правильно создавать потоки через Runnable и запускать их через start(). Возникает соблазн: «Отлично, теперь на каждый запрос пользователя я буду делать new Thread(task).start()!».

Это вторая фундаментальная ошибка, которая приведет к падению вашего приложения (Out Of Memory) при росте нагрузки.

Поток в Java — это тяжеловесный объект. Создание каждого нового потока требует времени на согласование с ОС. Кроме того, каждому потоку по умолчанию выделяется около 1 Мегабайта памяти под стек (Call Stack). Если к вам придет 5000 пользователей одновременно, создание 5000 потоков мгновенно съест 5 Гигабайт оперативной памяти только на пустые стеки, даже если задачи тривиальны.

Ручное создание потоков через new Thread() оправдано только для долгоживущих фоновых процессов. Для обработки множества коротких задач потоки нужно переиспользовать. Как именно это делается и почему пулы потоков спасают серверы от падения, мы разберем в следующих главах.

Жизненный цикл потока и управление его состоянием

Жизненный цикл потока и управление его состоянием

В прошлой главе мы выяснили, что создание потока в Java — это дорогая операция: операционная система выделяет ресурсы, а JVM резервирует около 1 МБ памяти под стек вызовов. Но что происходит с этими ресурсами, когда поток ждет ответа от базы данных или просто приостановлен? Сжигает ли он процессорное время впустую?

Чтобы писать эффективный многопоточный код, необходимо понимать, что поток не просто «работает» или «не работает». Он проживает сложный жизненный цикл, перемещаясь между различными состояниями под управлением планировщика операционной системы и виртуальной машины Java.

Шесть состояний потока в Java

Внутри класса Thread существует перечисление Thread.State, которое описывает шесть возможных состояний потока с точки зрения JVM.

  1. NEW — поток создан как объект в куче (вызван new Thread()), но нативный поток ОС еще не запущен.
  2. RUNNABLE — поток готов к выполнению или уже выполняется.
  3. BLOCKED — поток заблокирован в ожидании монитора (попытка войти в synchronized блок, занятый другим потоком).
  4. WAITING — поток бесконечно ждет действий от другого потока.
  5. TIMED_WAITING — поток ждет действий от другого потока в течение заданного времени.
  6. TERMINATED — поток успешно завершил метод run() или упал с необработанным исключением.

Важный инсайт: состояние RUNNABLE в Java скрывает в себе два разных состояния операционной системы — Ready (готов к выполнению, ждет очереди на ядро процессора) и Running (прямо сейчас выполняется на ядре). JVM не делает между ними различий. Если поток читает файл и физически заблокирован дисковой подсистемой ОС, для Java он все равно может числиться как RUNNABLE.

Управление выполнением: паузы и ожидания

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

Thread.sleep() — переход в TIMED_WAITING

Метод Thread.sleep(ms) приостанавливает выполнение текущего потока на указанное количество миллисекунд. Поток переходит в состояние TIMED_WAITING и гарантированно уступает процессор.

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

Thread.join() — ожидание завершения

Часто возникает ситуация, когда главному потоку нужно дождаться результатов работы фонового потока. Для этого используется метод join().

Если поток AA вызывает B.join(), то поток AA переходит в состояние WAITING (или TIMED_WAITING, если вызван join(ms)) и засыпает до тех пор, пока поток BB не перейдет в состояние TERMINATED.

Thread.yield() — вежливая уступка

Метод yield() — это подсказка планировщику ОС: «Я делаю не очень важную работу, если есть другие потоки в состоянии RUNNABLE, можешь отдать ядро им». При этом поток остается в состоянии RUNNABLE. В современном Java-программировании yield() практически не используется, так как его поведение сильно зависит от конкретной ОС и архитектуры процессора.

Как правильно остановить поток: Кооперативная отмена

Допустим, у нас есть фоновый поток, который бесконечно скачивает обновления. Как его остановить при закрытии приложения?

В ранних версиях Java существовал метод Thread.stop(). Он убивал поток мгновенно, в каком бы состоянии тот ни находился. Это приводило к катастрофам: открытые файлы не закрывались, соединения с базой данных повисали, а объекты оставались в наполовину измененном (неконсистентном) состоянии. Поэтому stop() был признан устаревшим (deprecated) и запрещен к использованию.

Вместо принудительного убийства в Java используется кооперативная отмена (cooperative cancellation) через механизм прерываний — interrupt().

Механизм прерывания (Interrupt)

Прерывание — это не команда «умри», это вежливая просьба: «Пожалуйста, заверши свою работу как можно скорее». У каждого потока есть внутренний boolean-флаг прерывания.

Когда вы вызываете thread.interrupt(), происходит одно из двух:

Сценарий 1: Поток активно работает (находится в RUNNABLE) Вызов interrupt() просто устанавливает внутренний флаг прерывания в true. Поток продолжит работу как ни в чем не бывало. Программист сам должен периодически проверять этот флаг внутри тяжелых циклов:

public class Worker implements Runnable {
    @Override
    public void run() {
        // Проверяем флаг прерывания на каждой итерации
        while (!Thread.currentThread().isInterrupted()) {
            doHeavyCalculations();
        }
        System.out.println("Поток получил сигнал и аккуратно завершается.");
        cleanUpResources();
    }
}

Сценарий 2: Поток спит или ждет (TIMED_WAITING / WAITING) Если поток заблокирован вызовом sleep(), join() или методами из пакета java.util.concurrent, он не может проверить флаг. В этом случае JVM немедленно будит поток и выбрасывает внутри него InterruptedException.

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

public void pauseExecution() {
    try {
        Thread.sleep(5000);
    } catch (InterruptedException e) {
        // Восстанавливаем флаг, так как sleep его сбросил!
        Thread.currentThread().interrupt();
        // Логируем или пробрасываем ошибку дальше
    }
}

Мы разобрали, как потоки живут, спят и корректно завершают работу. Однако до сих пор наши потоки существовали изолированно. Настоящие сложности начинаются, когда несколько потоков пытаются одновременно читать и изменять одни и те же данные. Именно там возникают состояния BLOCKED и потребность в синхронизации, к которым мы перейдем в следующей главе.

Первая безопасная передача данных между потоками

Первая безопасная передача данных между потоками

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

Анатомия памяти: Стек и Куча

Когда запускается новый поток, операционная система и JVM выделяют ему собственный стек (Stack). Стек — это личная, строго изолированная территория потока. В нём хранятся локальные переменные примитивных типов (int, double, boolean) и вызовы методов. Ни один другой поток в системе физически не может прочитать или изменить данные в чужом стеке.

Но объекты (экземпляры классов) никогда не живут в стеке. Они создаются в куче (Heap) — огромном общем пространстве памяти, доступном абсолютно всем потокам. В стеке потока хранится лишь ссылка на объект в куче.

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

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

Безопасные окна: start() и join()

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

JVM даёт нам строгую гарантию видимости памяти в двух точках жизненного цикла потока:

  1. При вызове Thread.start(). Всё, что главный поток записал в память до вызова start(), гарантированно будет видно новому потоку.
  2. При возврате из Thread.join(). Всё, что фоновый поток записал в память до своего завершения, гарантированно будет видно главному потоку после того, как join() успешно отработает.

Рассмотрим классический пример: нам нужно поручить фоновому потоку сложную математическую задачу, например, расчет факториала.

class FactorialTask implements Runnable {
    // Входные данные
    private final int number;
    // Результат
    private long result;

    public FactorialTask(int number) {
        this.number = number; // Передаем данные ДО старта
    }

    @Override
    public void run() {
        long calc = 1;
        for (int i = 1; i <= number; i++) {
            calc *= i;
        }
        this.result = calc; // Сохраняем результат ПЕРЕД завершением
    }

    public long getResult() {
        return result;
    }
}

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

FactorialTask task = new FactorialTask(10);
Thread worker = new Thread(task);

worker.start(); // Окно 1: worker гарантированно видит number = 10
worker.join();  // Окно 2: main гарантированно видит вычисленный result

System.out.println("Результат: " + task.getResult());

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

Проблема данных «в полёте»

Что делать, если данные нужно передать во время работы потока? Представьте сервер генерации отчетов. Фоновый поток постоянно формирует документы, а главный поток хочет на лету менять конфигурацию (например, формат вывода).

Если мы создадим класс ReportConfig с обычными геттерами и сеттерами:

class ReportConfig {
    private String format = "PDF";

    public void setFormat(String format) {
        this.format = format;
    }

    public String getFormat() {
        return format;
    }
}

Передача такого объекта в работающий поток — это бомба замедленного действия. Главный поток может вызвать setFormat("Excel") ровно в ту миллисекунду, когда фоновый поток вызывает getFormat() для формирования заголовка отчета. В результате часть отчета сгенерируется по правилам PDF, а часть — по правилам Excel. Произойдет нарушение целостности данных (Race Condition).

Ультимативная защита: Неизменяемость (Immutability)

Если проблема возникает из-за совместного изменения данных, самое элегантное решение — запретить изменения вообще. Объект, состояние которого невозможно изменить после создания, называется неизменяемым (Immutable).

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

Чтобы сделать класс неизменяемым в Java, необходимо выполнить четыре правила:

  1. Сделать класс final, чтобы его нельзя было расширить и переопределить методы.
  2. Сделать все поля private и final.
  3. Не предоставлять методы, меняющие состояние (сеттеры).
  4. Обеспечить глубокую неизменяемость: если поле является ссылкой на изменяемый объект (например, ArrayList или Date), класс не должен возвращать прямую ссылку на него.

Перепишем нашу конфигурацию отчетов по правилам Immutability:

public final class ReportConfig {
    private final String format;
    private final List<String> recipients;

    public ReportConfig(String format, List<String> recipients) {
        this.format = format;
        // Защитное копирование: создаем новую коллекцию,
        // чтобы вызывающий код не мог изменить список извне
        this.recipients = List.copyOf(recipients);
    }

    public String getFormat() {
        return format;
    }

    public List<String> getRecipients() {
        // List.copyOf возвращает неизменяемый список
        return recipients;
    }
}

Теперь, если главному потоку нужно изменить формат, он создает новый экземпляр ReportConfig("Excel", oldRecipients) и передает ссылку на него фоновому потоку. Старый объект останется в памяти (пока его не удалит сборщик мусора) и фоновый поток, который прямо сейчас генерирует PDF, спокойно завершит свою работу со старой, целостной конфигурацией.

Использование start(), join() и неизменяемых объектов — это фундамент безопасного многопоточного программирования. Эти подходы позволяют избежать 90% ошибок, связанных с гонкой данных, не прибегая к тяжеловесным блокировкам. Однако бывают ситуации, когда потокам всё же необходимо совместно обновлять единый счетчик или общую структуру данных. Для этого потребуются механизмы синхронизации.

Ограничение доступа к состоянию: ThreadLocal на практике

Ограничение доступа к состоянию: ThreadLocal на практике

Мы уже знаем, что объекты в куче (Heap) доступны всем потокам, и самый надежный способ безопасно передать данные — сделать их неизменяемыми (Immutable). Но как быть, если объект принципиально имеет изменяемое состояние, а создавать его копию для каждой операции слишком дорого?

Классический пример — SimpleDateFormat в Java. Он не является потокобезопасным: если несколько потоков одновременно вызовут у одного экземпляра метод format(), внутреннее состояние объекта (календарь) будет перезаписано непредсказуемым образом, и на выходе мы получим неверные даты или исключения.

Создавать new SimpleDateFormat() при каждом вызове метода — значит создавать огромную нагрузку на сборщик мусора (Garbage Collector). Нам нужно поведение локальной переменной (изоляция), но с временем жизни, выходящим за рамки одного вызова метода.

Этот подход называется ограничением потоком (Thread Confinement): если объект не является потокобезопасным, мы просто не делимся им. Мы выдаем каждому потоку его личный, эксклюзивный экземпляр.

В Java этот паттерн реализуется с помощью класса ThreadLocal.

Базовый синтаксис ThreadLocal

ThreadLocal работает как контейнер, который предоставляет каждому обращающемуся к нему потоку собственную, независимо инициализированную копию значения.

public class DateFormatter {
    // Инициализация значения для каждого нового потока
    private static final ThreadLocal<SimpleDateFormat> dateFormat =
        ThreadLocal.withInitial(() -> new SimpleDateFormat("yyyy-MM-dd"));

    public String format(Date date) {
        // get() вернет экземпляр, привязанный к текущему потоку
        return dateFormat.get().format(date);
    }
}

Когда поток впервые вызывает dateFormat.get(), ThreadLocal выполняет лямбду из withInitial(), создает новый SimpleDateFormat и сохраняет его специально для этого потока. При последующих вызовах get() тот же поток получит тот же самый экземпляр. Другой поток получит свой собственный экземпляр. Никакой гонки данных нет, так как нет разделяемого состояния.

Как это работает под капотом (и почему это не Map)

Интуитивно кажется, что внутри ThreadLocal находится глобальная потокобезопасная коллекция (например, ConcurrentHashMap), где ключом выступает Thread.currentThread(), а значением — сохраненный объект.

Ранние версии Java так и работали, но этот подход имеет фатальный недостаток: глобальная мапа становится узким местом (bottleneck) при синхронизации. Если тысячи потоков одновременно попытаются получить свои значения, они выстроятся в очередь к одной структуре данных.

Современная реализация вывернута наизнанку.

Глобального хранилища не существует. Вместо этого сам класс Thread имеет скрытое поле типа ThreadLocal.ThreadLocalMap.

Когда вы вызываете threadLocal.get(), под капотом происходит следующее:

  1. Берется текущий поток: Thread t = Thread.currentThread();
  2. У этого потока запрашивается его личная мапа: ThreadLocalMap map = t.threadLocals;
  3. В этой мапе ищется значение, где ключом выступает сам объект ThreadLocal (точнее, ссылка на него).

Данные хранятся не в объекте ThreadLocal, а внутри самого объекта Thread. ThreadLocal — это лишь ключ для доступа к слоту в памяти текущего потока.

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

Темная сторона: утечки памяти и переиспользование потоков

Поскольку данные лежат внутри объекта Thread, они будут жить в памяти ровно столько, сколько живет сам поток. Если мы создали поток через new Thread(), выполнили задачу и поток завершился — сборщик мусора очистит и объект Thread, и его ThreadLocalMap, и все сохраненные там значения.

Однако в реальных enterprise-приложениях (например, на базе Spring или Tomcat) потоки почти никогда не создаются на один раз. Они живут в пулах потоков (Thread Pools). Поток берется из пула, обрабатывает HTTP-запрос пользователя, а затем возвращается в пул, чтобы обработать следующий запрос (возможно, уже от другого пользователя).

И здесь кроется главная опасность ThreadLocal.

Допустим, мы сохранили в ThreadLocal ID текущего пользователя, чтобы не передавать его через параметры всех методов бизнес-логики:

public class UserContext {
    public static final ThreadLocal<String> currentUser = new ThreadLocal<>();
}

// Где-то в начале обработки запроса
UserContext.currentUser.set(request.getUserId());

Если после завершения обработки запроса мы не очистим это значение, поток вернется в пул вместе с заполненным ThreadLocalMap. Когда этот же поток возьмут для обработки запроса другого пользователя, метод UserContext.currentUser.get() вернет старый ID. Это приведет к критической уязвимости: один пользователь получит доступ к данным другого.

Кроме того, если в ThreadLocal сохраняются тяжелые объекты (например, кэши или загрузчики классов), а потоки живут неделями, это приводит к классической утечке памяти (Memory Leak) — объекты накапливаются в пуле потоков и никогда не собираются Garbage Collector'ом.

Правильный паттерн использования

Чтобы избежать проблем при переиспользовании потоков, работа с ThreadLocal всегда должна оборачиваться в блок try-finally с обязательным вызовом метода remove().

public void processRequest(Request request) {
    UserContext.currentUser.set(request.getUserId());
    try {
        // Вся бизнес-логика выполняется здесь
        // Любой метод внутри может вызвать UserContext.currentUser.get()
        handle(request);
    } finally {
        // Гарантированное удаление значения из ThreadLocalMap текущего потока,
        // даже если handle(request) выбросил исключение
        UserContext.currentUser.remove();
    }
}

Метод remove() удаляет запись (ключ-значение) из ThreadLocalMap текущего потока. Теперь, когда поток вернется в пул, он будет абсолютно "чистым" и не сохранит ссылок на данные обработанного запроса.

Проблема гонки данных (Race Condition) и критические секции

Проблема гонки данных (Race Condition) и критические секции

До сих пор мы избегали совместного изменения данных. Мы прятали состояние внутри стека, изолировали его через ThreadLocal или «замораживали» с помощью Immutable-объектов. Но в реальных высоконагруженных системах потокам необходимо общаться и менять общее состояние: обновлять кэши, увеличивать счетчики метрик, списывать средства со счетов.

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

Иллюзия атомарности

Представьте простейшую задачу: нам нужен глобальный счетчик посещений.

public class HitCounter {
    private int count = 0;

    public void increment() {
        count++;
    }

    public int getCount() {
        return count;
    }
}

Если мы запустим два потока, и каждый вызовет increment() по 10 000 раз, мы ожидаем увидеть в результате 20 000. Но на практике мы получим 15 432, 18 991 или любое другое случайное число меньше ожидаемого. Куда пропадают инкременты?

Проблема кроется в операции count++. В исходном коде на Java это выглядит как одно неделимое действие. Однако процессор не понимает Java. На уровне машинных инструкций инкремент распадается на три независимых шага:

  1. LOAD: Прочитать текущее значение count из оперативной памяти в регистр процессора.
  2. ADD: Увеличить значение в регистре на 1.
  3. STORE: Записать новое значение из регистра обратно в память.

Операция называется атомарной (от греч. atomos — неделимый), если она выполняется целиком или не выполняется вовсе, и никакой другой поток не может увидеть ее промежуточное состояние. Выражение count++ не является атомарным.

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

Ситуация, при которой результат выполнения программы зависит от того, в каком порядке планировщик ОС переключает потоки, называется состоянием гонки или гонкой данных (Race Condition).

Два главных паттерна Race Condition

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

1. Read-Modify-Write (Прочитай-Измени-Запиши)

Это именно тот случай, который мы разобрали со счетчиком count++. Поток читает старое состояние, вычисляет на его основе новое и записывает обратно. Если между чтением и записью другой поток успеет изменить исходные данные, первый поток запишет результат, основанный на уже устаревшей информации. Возникает эффект «потерянного обновления» (Lost Update).

2. Check-Then-Act (Проверь-Затем-Действуй)

Второй паттерн возникает, когда поток принимает решение на основе какого-то условия, но к моменту выполнения действия условие уже перестает быть истинным.

Классический пример — ленивая инициализация (Lazy Initialization) без синхронизации:

public class ExpensiveResourceFactory {
    private Resource instance = null;

    public Resource getInstance() {
        if (instance == null) { // Check
            instance = new Resource(); // Act
        }
        return instance;
    }
}

Что произойдет, если Поток А и Поток Б вызовут getInstance() одновременно?

  1. Поток А проверяет instance == null. Это true.
  2. Планировщик ставит Поток А на паузу.
  3. Поток Б проверяет instance == null. Переменная все еще null, потому что Поток А не успел создать объект. Это true.
  4. Поток Б создает new Resource() и записывает ссылку в instance.
  5. Поток А просыпается и продолжает работу с того места, где остановился. Он тоже создает new Resource() и перезаписывает ссылку в instance.

Результат: создано два тяжелых объекта вместо одного, первый объект потерян (уйдет в Garbage Collector), а логика синглтона сломана.

Критическая секция и взаимное исключение

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

Участок кода, в котором происходит доступ к разделяемому изменяемому состоянию и который не должен выполняться более чем одним потоком одновременно, называется критической секцией (Critical Section).

Свойство системы, гарантирующее, что в критической секции в любой момент времени находится максимум один поток, называется взаимным исключением (Mutual Exclusion или Mutex).

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

  1. Первый поток подходит к секции, видит зеленый свет, включает красный для остальных и заходит внутрь.
  2. Пока он внутри, любые другие потоки, дошедшие до этого места, вынуждены ждать (блокироваться).
  3. Когда первый поток выходит, он снова включает зеленый свет, позволяя зайти следующему.

Важно понимать: критическая секция — это не свойство самих данных в памяти. Это договоренность на уровне логики программы. Мы берем набор инструкций (например, чтение count, инкремент, запись) и оборачиваем их в защитный барьер.

В Java основным встроенным механизмом для создания таких светофоров и выделения критических секций является ключевое слово synchronized и концепция мониторов. Именно их внутреннее устройство и правила применения мы разберем в следующем шаге.

Ключевое слово synchronized и мониторы объектов

Ключевое слово synchronized и мониторы объектов

Мы знаем, что неатомарные операции и паттерны вроде Check-Then-Act приводят к состоянию гонки, если несколько потоков одновременно обращаются к общим данным. Чтобы этого избежать, нам нужен механизм взаимного исключения — способ создать критическую секцию, в которой единовременно может находиться только один поток. В Java самым базовым и встроенным прямо в язык инструментом для этого является ключевое слово synchronized.

Синтаксис взаимного исключения

В Java есть два способа объявить критическую секцию с помощью synchronized: на уровне всего метода и на уровне отдельного блока кода.

Синхронизированный метод объявляется добавлением ключевого слова в сигнатуру:

public synchronized void increment() {
    count++; // Критическая секция
}

Синхронизированный блок требует явного указания объекта, который будет выступать в роли «замка»:

public void increment() {
    synchronized(this) {
        count++; // Критическая секция
    }
}

Технически, первый вариант — это просто синтаксический сахар для второго. Когда вы объявляете нестатический метод как synchronized, Java автоматически использует текущий экземпляр объекта (this) в качестве замка. Если метод статический, замком выступает объект класса (MyClass.class).

Но что именно означает «использовать объект в качестве замка»?

Анатомия блокировки: что такое Монитор

В Java не существует отдельного базового класса «Мьютекс», который нужно было бы создавать для каждой блокировки. Создатели языка заложили механизм синхронизации в фундамент: абсолютно любой объект в Java может выступать в роли блокировки.

Это возможно благодаря тому, что за каждым объектом закреплена скрытая сущность, которая называется монитором (Monitor).

Монитор — это механизм контроля доступа. У него есть владелец (поток, который его захватил) и очередь ожидания. Когда поток встречает ключевое слово synchronized(obj), он обращается к монитору объекта obj.

Связь между объектом и его монитором хранится в заголовке объекта в памяти. Каждый объект в куче (Heap) начинается с метаданных — Object Header. Часть этого заголовка называется Mark Word. Именно в Mark Word JVM записывает информацию о состоянии блокировки. Когда объект не заблокирован, Mark Word хранит хеш-код и флаги сборщика мусора. Но как только поток захватывает объект через synchronized, Mark Word переписывается и начинает указывать на структуру монитора в памяти JVM.

Как потоки выстраиваются в очередь

Рассмотрим динамику захвата монитора. Допустим, у нас есть объект account и три потока, которые одновременно пытаются выполнить synchronized(account).

  1. Захват: Поток А первым добирается до инструкции synchronized. JVM проверяет монитор объекта account. Монитор свободен. Поток А становится владельцем монитора и заходит в критическую секцию.
  2. Блокировка: Поток B добирается до synchronized(account). JVM видит, что монитор занят Потоком А. Поток B переводится ОС в состояние BLOCKED (заблокирован) и помещается в очередь ожидания (Entry Set) этого монитора. Он не потребляет процессорное время, он просто спит. То же самое происходит с Потоком C.
  3. Освобождение: Поток А завершает выполнение блока synchronized (или выбрасывает исключение — блокировка снимается гарантированно). Монитор освобождается.
  4. Пробуждение: JVM выбирает один из потоков в очереди ожидания (например, Поток B), переводит его обратно в состояние RUNNABLE, и тот захватывает монитор.

Состояние BLOCKED — это специфическое состояние жизненного цикла потока в Java, которое возникает только при ожидании захвата монитора перед входом в synchronized блок или метод.

Реентерабельность (Повторная входимость)

Представьте ситуацию: один synchronized метод внутри класса вызывает другой synchronized метод этого же класса.

public synchronized void process() {
    stepOne();
    stepTwo();
}

public synchronized void stepOne() {
    // логика
}

Когда поток вызывает process(), он захватывает монитор объекта this. Внутри он вызывает stepOne(), который тоже требует монитор объекта this. Если бы монитор был простой защелкой, поток заблокировал бы сам себя навсегда, ожидая освобождения монитора, который он же и держит.

Этого не происходит, потому что мониторы в Java реентерабельны (reentrant).

Внутри структуры монитора хранятся два важных поля:

  1. Идентификатор потока-владельца.
  2. Счетчик вхождений (reentrancy count).

Когда поток впервые захватывает свободный монитор, счетчик становится равен 1, а поток записывается как владелец. Если этот же поток снова встречает блок synchronized по тому же объекту, JVM видит совпадение владельца и просто увеличивает счетчик до 2. При выходе из вложенного блока счетчик уменьшается. Монитор считается полностью свободным и отдается другим потокам только тогда, когда счетчик опускается до 0.

Гранулярность блокировок

Использование synchronized на уровне метода — это просто, но часто ведет к проблемам с производительностью.

Представьте класс, который собирает статистику веб-сервера: подсчитывает количество успешных запросов и количество ошибок.

public class ServerStats {
    private int successCount = 0;
    private int errorCount = 0;

    public synchronized void addSuccess() {
        successCount++;
    }

    public synchronized void addError() {
        errorCount++;
    }
}

Поскольку оба метода synchronized, они оба используют один и тот же монитор — this. Если Поток А вызывает addSuccess(), он блокирует весь объект. Если в этот момент Поток B попытается вызвать addError(), он будет заблокирован и переведен в BLOCKED.

Но переменные successCount и errorCount независимы! Нет никакой логической причины запрещать одновременное обновление успехов и ошибок. Используя synchronized на методах, мы сделали гранулярность блокировки слишком крупной, создав искусственное узкое место (bottleneck).

Правильный подход — использовать отдельные объекты-замки для независимых данных:

public class ServerStats {
    private int successCount = 0;
    private int errorCount = 0;

    // Специальные объекты-замки
    private final Object successLock = new Object();
    private final Object errorLock = new Object();

    public void addSuccess() {
        synchronized(successLock) {
            successCount++;
        }
    }

    public void addError() {
        synchronized(errorLock) {
            errorCount++;
        }
    }
}

Теперь addSuccess() и addError() используют разные мониторы. Потоки, обновляющие разные метрики, больше не блокируют друг друга, при этом каждый счетчик надежно защищен от состояния гонки.

Ключевое слово synchronized — это мощный и надежный инструмент. Однако его простота обманчива. Захват монитора гарантирует взаимное исключение, но не решает проблему координации: что делать, если поток вошел в критическую секцию, но данные для работы еще не готовы? Для этого мониторы предоставляют дополнительный механизм управления ожиданием, который мы разберем далее.

Взаимодействие потоков через wait, notify и notifyAll

Взаимодействие потоков через wait, notify и notifyAll

Представьте, что поток захватил монитор объекта, зашел в критическую секцию и обнаружил, что нужные ему данные еще не готовы. Что ему делать? Если он запустит бесконечный цикл проверки условия (активное ожидание или busy-waiting), он сожжет процессорное время впустую. Если он просто уснет через метод sleep(), он заблокирует монитор, и другие потоки не смогут войти в критическую секцию, чтобы подготовить эти самые данные. Возникнет тупик.

Нам нужен механизм координации: поток должен уметь сказать «я подожду, пока условие не выполнится, а пока забирайте монитор». Именно для этого в Java существуют методы wait(), notify() и notifyAll().

Две комнаты монитора: Entry Set и Wait Set

Ранее мы выяснили, что каждый объект в Java имеет ассоциированный с ним монитор. До сих пор мы рассматривали монитор как комнату, в которой может находиться только один поток-владелец, а остальные толпятся снаружи в очереди. Эта очередь называется Entry Set (набор входа), и потоки в ней находятся в состоянии BLOCKED.

Но внутри монитора есть еще одна скрытая зона — Wait Set (набор ожидания).

Wait Set — это специальная очередь для потоков, которые добровольно временно отказались от владения монитором, потому что ждут наступления определенного логического условия. Потоки в этой зоне находятся в состоянии WAITING (или TIMED_WAITING, если вызван метод с таймаутом).

Механика метода wait()

Методы wait(), notify() и notifyAll() определены прямо в классе Object, а не в классе Thread. Это логично, потому что они управляют монитором конкретного объекта, а не самим потоком напрямую.

Главное правило: вызывать эти методы можно только изнутри критической секции, то есть поток уже должен владеть монитором этого объекта. Если вызвать wait() без блока synchronized, JVM немедленно выбросит IllegalMonitorStateException.

Когда поток вызывает object.wait(), под капотом происходит следующая последовательность:

  1. Поток освобождает монитор объекта (сбрасывает счетчик реентерабельности до нуля).
  2. Поток помещается в Wait Set этого монитора.
  3. Состояние потока меняется на WAITING, и планировщик ОС перестает выделять ему процессорное время.
  4. Монитор становится свободен, и один из потоков из Entry Set может его захватить.

Пробуждение: notify() и notifyAll()

Поток может находиться в Wait Set бесконечно долго. Чтобы вернуть его к жизни, другой поток, захвативший этот же монитор, должен изменить данные (выполнить условие) и подать сигнал.

Вызов object.notify() берет один случайный поток из Wait Set и переводит его в Entry Set. Вызов object.notifyAll() переводит все потоки из Wait Set в Entry Set.

Важнейший нюанс: вызов notify() не означает, что разбуженный поток тут же начнет выполняться. Метод notify() не освобождает монитор! Разбуженный поток просто переходит из состояния WAITING в состояние BLOCKED и встает в общую очередь Entry Set. Он продолжит выполнение (выйдет из метода wait()) только тогда, когда текущий владелец монитора выйдет из блока synchronized, и разбуженному потоку удастся выиграть конкуренцию за захват блокировки.

В 99% случаев профессионалы используют notifyAll(). Использование одиночного notify() опасно проблемой «потерянного сигнала» (lost wakeup): если в Wait Set ждут потоки разных типов (например, читатели и писатели), случайный сигнал может разбудить поток, чье условие еще не выполнено, а нужный поток так и останется спать навсегда.

Золотое правило: всегда используйте цикл while

Рассмотрим классический код ожидания условия:

synchronized (lock) {
    while (!conditionIsMet) {
        lock.wait();
    }
    // Выполнение полезной работы
}

Почему используется цикл while, а не условный оператор if? На это есть две фундаментальные причины.

Первая причина — изменение состояния (State Change). Представьте, что поток А ждет 데이터를, а поток Б вызывает notifyAll(). Поток А просыпается и переходит в Entry Set. Но пока поток А ждал своей очереди на захват монитора, шустрый поток В успел захватить монитор первым, изменить данные и сделать условие снова ложным. Когда поток А наконец войдет в секцию, данные уже невалидны. Если бы он использовал if, он бы пошел выполнять работу с неверными данными. Цикл while заставит его снова проверить условие и, если нужно, опять уснуть.

Вторая причина — ложное пробуждение (Spurious Wakeup). Это специфика работы планировщиков потоков на уровне операционных систем (POSIX Threads). ОС имеет право разбудить поток, находящийся в состоянии ожидания, даже если никто не вызывал метод notify(). Это происходит редко, но спецификация Java прямо говорит: поток может проснуться без причины. Обертка в while надежно защищает от этой аномалии.

Практика: паттерн Producer-Consumer

Соберем все вместе на классической задаче «Производитель-Потребитель». У нас есть буфер ограниченного размера. Производители кладут туда сообщения, потребители — забирают. Если буфер полон, производитель должен ждать. Если пуст — ждет потребитель.

public class MessageBroker {
    private String message;
    private boolean hasMessage = false;

    // Метод для Производителя
    public synchronized void produce(String msg) throws InterruptedException {
        while (hasMessage) {
            wait(); // Ждем, пока буфер не освободится
        }
        this.message = msg;
        this.hasMessage = true;
        System.out.println("Произвел: " + msg);
        notifyAll(); // Будим Потребителей
    }

    // Метод для Потребителя
    public synchronized String consume() throws InterruptedException {
        while (!hasMessage) {
            wait(); // Ждем, пока не появится сообщение
        }
        String msg = this.message;
        this.hasMessage = false;
        System.out.println("Потребил: " + msg);
        notifyAll(); // Будим Производителей
        return msg;
    }
}

В этом коде монитором выступает сам экземпляр MessageBroker (так как методы synchronized). И производители, и потребители делят один и тот же Wait Set. Когда потребитель забирает сообщение, он вызывает notifyAll(). Это будит всех: и других потребителей, и производителей. Потребители проснутся, проверят цикл while(!hasMessage), увидят, что сообщений нет, и снова уснут. А производитель увидит, что место освободилось, и положит новое сообщение.

Прямое использование wait/notify требует высокой дисциплины: легко забыть цикл while, легко перепутать мониторы, легко получить зависание из-за использования обычного notify. Поэтому в современном Java-коде эти низкоуровневые примитивы используются редко прикладными программистами.

Проблема взаимной блокировки (Deadlock) и способы её предотвращения

Проблема взаимной блокировки (Deadlock) и способы её предотвращения

Представьте классическую задачу: перевод денег между двумя банковскими счетами. У нас есть класс Account и метод transfer, который должен списать сумму с одного счета и зачислить на другой. Чтобы избежать состояния гонки (Race Condition), мы защищаем оба счета мониторами:

public void transfer(Account from, Account to, int amount) {
    synchronized (from) {
        synchronized (to) {
            if (from.getBalance() >= amount) {
                from.withdraw(amount);
                to.deposit(amount);
            }
        }
    }
}

Код выглядит абсолютно логичным и потокобезопасным. Но если Алиса решит перевести деньги Бобу, и ровно в ту же миллисекунду Боб решит перевести деньги Алисе, приложение застынет навсегда. Не будет выброшено исключений, процессор не покажет 100% нагрузки. Потоки просто уснут, ожидая друг друга. Это классический пример взаимной блокировки — Deadlock.

Анатомия катастрофы

Разберем по шагам, что происходит в памяти при встречном переводе.

Поток 1 (перевод от Алисы к Бобу):

  1. Захватывает монитор объекта Alice.
  2. Пытается захватить монитор объекта Bob.

Поток 2 (перевод от Боба к Алисе):

  1. Захватывает монитор объекта Bob.
  2. Пытается захватить монитор объекта Alice.

Если планировщик ОС переключит контекст ровно между первым и вторым шагами, возникнет патовая ситуация. Поток 1 держит Alice и ждет Bob. Поток 2 держит Bob и ждет Alice. Ни один из них не может продолжить работу, чтобы освободить уже захваченный ресурс.

Встроенные мониторы Java (synchronized) не имеют механизма таймаутов или принудительного отзыва блокировки. Если поток вошел в состояние ожидания захвата монитора (состояние BLOCKED), он будет ждать вечно.

Четыре условия Коффмана

В 1971 году информатик Эдвард Коффман формализовал теорию взаимных блокировок. Он доказал, что дедлок возможен тогда и только тогда, когда в системе одновременно выполняются четыре условия.

Эдвард Коффман

  1. Взаимное исключение (Mutual Exclusion). Ресурс может удерживаться только одним потоком одновременно. В нашем случае это обеспечивается самим смыслом блока synchronized.
  2. Удержание и ожидание (Hold and Wait). Поток, уже удерживающий как минимум один ресурс, запрашивает дополнительные ресурсы, которые удерживаются другими потоками.
  3. Отсутствие вытеснения (No Preemption). Ресурс не может быть принудительно забран у удерживающего его потока; поток должен освободить его добровольно после завершения работы.
  4. Круговое ожидание (Circular Wait). Существует замкнутая цепь из двух и более потоков, в которой каждый поток ждет ресурс, удерживаемый следующим потоком в цепи.

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

Посмотрим на арсенал Java через призму этих условий. Можем ли мы отказаться от взаимного исключения? Нет, иначе мы получим Race Condition при изменении баланса. Можем ли мы реализовать вытеснение (забрать монитор у другого потока)? Нет, JVM не позволяет прерывать потоки, находящиеся в synchronized блоке.

У нас остается два пути: избавиться от «удержания и ожидания» или разорвать «круговое ожидание».

Разрыв круга: глобальное упорядочивание блокировок

Самый надежный и распространенный способ предотвращения дедлоков при использовании synchronized — нарушение условия кругового ожидания. Это достигается техникой Lock Ordering (упорядочивание блокировок).

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

Вернемся к банковским счетам. У каждого счета в базе данных есть уникальный идентификатор (ID). Мы можем ввести правило: всегда сначала захватывать монитор счета с меньшим ID, а затем — с большим.

public void transferSafe(Account a, Account b, int amount) {
    Account firstLock = a.getId() < b.getId() ? a : b;
    Account secondLock = a.getId() < b.getId() ? b : a;

    synchronized (firstLock) {
        synchronized (secondLock) {
            if (a.getBalance() >= amount) {
                a.withdraw(amount);
                b.deposit(amount);
            }
        }
    }
}

Что теперь произойдет при встречном переводе? Пусть ID Алисы = 100, а ID Боба = 200.

  • Поток 1 (Алиса \to Боб) вычислит: firstLock = Алиса, secondLock = Боб.
  • Поток 2 (Боб \to Алиса) вычислит: firstLock = Алиса, secondLock = Боб.

Оба потока попытаются сначала захватить монитор Алисы. Тот поток, который успеет первым, продолжит работу и захватит монитор Боба. Второй поток будет мирно ждать на мониторе Алисы. Дедлок невозможен.

Проблема одинаковых ключей (Tie-Breaking)

Что делать, если у объектов нет естественного уникального идентификатора? В Java можно использовать адрес объекта в памяти, который возвращает метод System.identityHashCode(Object).

Однако существует микроскопическая вероятность коллизии — два разных объекта могут вернуть одинаковый хэш-код. Если firstLock и secondLock окажутся неразличимы для нашей логики сортировки, мы снова рискуем получить разный порядок захвата в разных потоках.

Для разрешения таких ситуаций вводится третий, глобальный замок — Tie-Breaking Lock (блокировка разрешения ничьей).

private static final Object TIE_LOCK = new Object();

public void transferWithHash(Account a, Account b, int amount) {
    int hashA = System.identityHashCode(a);
    int hashB = System.identityHashCode(b);

    if (hashA < hashB) {
        synchronized (a) { synchronized (b) { doTransfer(a, b, amount); } }
    } else if (hashA > hashB) {
        synchronized (b) { synchronized (a) { doTransfer(a, b, amount); } }
    } else {
        // Коллизия хэш-кодов! Используем глобальный замок.
        synchronized (TIE_LOCK) {
            synchronized (a) { synchronized (b) { doTransfer(a, b, amount); } }
        }
    }
}

Глобальный замок TIE_LOCK становится узким местом (bottleneck) для производительности, так как сериализует все операции с коллизиями. Но поскольку коллизии identityHashCode крайне редки, на практике это не влияет на пропускную способность системы, зато гарантирует 100% защиту от дедлока.

Ограничения встроенных мониторов

Упорядочивание блокировок работает отлично, когда мы знаем все необходимые ресурсы заранее. Но что, если список ресурсов формируется динамически в процессе работы, или архитектура приложения не позволяет выстроить единую иерархию объектов?

В таких случаях нам нужно нарушить другое условие Коффмана — «Удержание и ожидание» (Hold and Wait). Поток должен попытаться захватить вторую блокировку, и, если она занята, освободить первую, немного подождать и попытаться снова.

Сделать это с помощью блоков synchronized невозможно, так как при входе в блок поток либо захватывает монитор, либо блокируется навсегда. Для реализации таких сложных сценариев в Java существуют явные блокировки из пакета java.util.concurrent.locks, которые позволяют запрашивать доступ с таймаутом.

Явные блокировки Lock и ReentrantLock

Явные блокировки Lock и ReentrantLock

Встроенные мониторы и ключевое слово synchronized отлично справляются с базовой защитой критических секций. Но у них есть фатальный недостаток: если поток попытался войти в synchronized-блок, а монитор уже занят, поток переходит в состояние ожидания навсегда. Вы не можете сказать ему: «подожди 5 секунд, и если не получится — отмени операцию». Вы не можете прервать это ожидание.

В прошлой главе мы разбирали взаимную блокировку (дедлок) и упорядочивали мониторы по ID, чтобы разорвать круговое ожидание. Но что, если мы могли бы просто нарушить другое условие Коффмана — «Удержание и ожидание»? Что, если поток мог бы попытаться захватить второй замок, потерпеть неудачу, отпустить первый замок и попробовать всё заново?

Для таких сценариев в Java существует пакет java.util.concurrent.locks и интерфейс Lock.

Отказ от неявности: базовый синтаксис

Интерфейс Lock предоставляет те же гарантии видимости памяти и взаимного исключения, что и synchronized. Самая популярная его реализация — класс ReentrantLock. Как следует из названия, он обладает свойством реентерабельности: поток может многократно захватывать блокировку, которой уже владеет (точно так же, как мы обсуждали это для мониторов).

Главное отличие Lock — управление блокировкой становится явным. Вы сами вызываете методы захвата и освобождения.

Lock lock = new ReentrantLock();

lock.lock(); // Поток блокируется здесь, если замок занят
try {
    // Критическая секция
    processSharedResource();
} finally {
    lock.unlock(); // Обязательное освобождение
}

Ключевое правило использования явных блокировок: метод unlock() всегда должен вызываться в блоке finally. Если внутри критической секции вылетит исключение (например, NullPointerException), а unlock() останется в основном блоке try, замок никогда не будет освобожден. Все остальные потоки, ожидающие этот Lock, зависнут навсегда.

tryLock: элегантное решение проблемы дедлоков

Главное оружие Lock против взаимных блокировок — метод tryLock(). В отличие от обычного lock(), он не погружает поток в ожидание. Если замок свободен, метод захватывает его и возвращает true. Если занят — мгновенно возвращает false, позволяя потоку продолжить выполнение и заняться чем-то другим.

Существует также перегруженная версия tryLock(long time, TimeUnit unit), которая ждет указанное время, прежде чем сдаться.

Вернемся к примеру с переводом денег между двумя банковскими счетами из прошлой главы. Вместо сложной сортировки счетов по ID (Lock Ordering), мы можем использовать tryLock(). Алгоритм будет таким:

  1. Пытаемся захватить замок первого счета.
  2. Если удалось — пытаемся захватить замок второго счета.
  3. Если второй замок занят, мы обязательно отпускаем первый и уходим на небольшую паузу.
  4. Повторяем попытку.
public void transfer(Account from, Account to, BigDecimal amount) {
    while (true) {
        if (from.getLock().tryLock()) {
            try {
                if (to.getLock().tryLock()) {
                    try {
                        // Оба замка захвачены, выполняем перевод
                        from.withdraw(amount);
                        to.deposit(amount);
                        return; // Успешный выход из цикла
                    } finally {
                        to.getLock().unlock();
                    }
                }
            } finally {
                from.getLock().unlock();
            }
        }

        // Если мы здесь, значит не удалось захватить оба замка.
        // Спим случайное время, чтобы избежать Livelock (когда потоки
        // синхронно захватывают первые замки и отпускают их)
        Thread.sleep(Random.nextInt(10));
    }
}

Этот подход нарушает условие «Удержание и ожидание»: поток никогда не ждет второй ресурс, удерживая первый.

Справедливость (Fairness) и проблема голодания

Встроенный synchronized не дает никаких гарантий о том, в каком порядке потоки получат доступ к монитору. Если 10 потоков ждут освобождения критической секции, JVM может выбрать любой из них. Часто это приводит к тому, что один «невезучий» поток может ждать очень долго, пока другие постоянно захватывают монитор. Это явление называется голоданием потока (Thread Starvation).

При создании ReentrantLock вы можете передать в конструктор boolean флаг, определяющий политику справедливости:

  • new ReentrantLock(false) (по умолчанию) — нечестная блокировка.
  • new ReentrantLock(true) — честная блокировка (Fair Lock).

В честной блокировке потоки выстраиваются в строгую FIFO-очередь (First In, First Out). Замок всегда достается тому потоку, который ждал дольше всех.

Возникает вопрос: если честная блокировка защищает от голодания, почему по умолчанию ReentrantLock создается нечестным?

Дело в производительности. Возобновление приостановленного потока (Context Switch) — тяжелая операция для операционной системы. Представьте ситуацию: поток AA только что отпустил замок. В очереди давно ждет поток BB, и ОС начинает его «будить». В эту же микросекунду к замку подбегает совершенно новый, уже активный поток CC.

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

Объекты Condition: точечное управление ожиданием

В главе про взаимодействие потоков мы выяснили, что у каждого объекта-монитора есть только один Wait Set. Из-за этого вызов notifyAll() будит вообще всех ожидающих, даже если условие выполнилось только для одного типа потоков.

Например, в ограниченном буфере (Producer-Consumer) производители ждут, пока появится свободное место, а потребители — пока появятся данные. При использовании wait/notifyAll освобождение одного места будит не только производителей, но и других потребителей, которые просыпаются впустую (ведь данных больше не стало).

Lock решает эту проблему через интерфейс Condition. Один Lock может иметь сколько угодно объектов Condition, каждый со своей собственной очередью ожидания.

Lock lock = new ReentrantLock();
// Создаем два независимых условия, привязанных к одному замку
Condition notFull = lock.newCondition();
Condition notEmpty = lock.newCondition();

Методы Condition работают аналогично методам класса Object, но называются иначе, чтобы избежать путаницы:

  • Вместо wait() используется await().
  • Вместо notify() используется signal().
  • Вместо notifyAll() используется signalAll().

Теперь при добавлении элемента в буфер производитель вызывает notEmpty.signal(). Это разбудит только одного потребителя из очереди notEmpty. Производители, ожидающие в очереди notFull, останутся спать. Это кардинально снижает количество ложных пробуждений и экономит ресурсы процессора в высоконагруженных системах.

Важно помнить: как и wait(), метод await() можно вызывать только тогда, когда поток уже владеет замком (находится внутри lock.lock()). При вызове await() поток атомарно отпускает замок и засыпает, а при пробуждении — заново борется за его захват.

Координация фаз работы: CountDownLatch и CyclicBarrier

Координация фаз работы: CountDownLatch и CyclicBarrier

Представьте, что ваш сервер при запуске должен загрузить данные из трех независимых источников: профили пользователей из базы данных, настройки из конфигурационного файла и горячие данные из Redis. Сервер не должен принимать входящие HTTP-запросы, пока все три кэша не будут готовы.

Как это реализовать? Можно создать три потока, сохранить ссылки на них и у каждого вызвать метод join(). Но что, если загрузка кэшей — это лишь часть работы пула потоков, и сами потоки не завершаются после этой задачи? Использовать wait() и notifyAll()? Код быстро обрастет флагами состояний, блоками synchronized и циклами проверки, превратившись в хрупкую конструкцию.

Для решения высокоуровневых задач координации в пакете java.util.concurrent предусмотрены готовые синхронизаторы. Они инкапсулируют внутри себя сложную логику работы с Lock и Condition, предоставляя простой и надежный API.

CountDownLatch: одноразовая защелка

CountDownLatch (защелка с обратным отсчетом) позволяет одному или нескольким потокам ожидать завершения набора операций, выполняемых другими потоками.

Механика работы предельно проста:

  1. При создании объекта задается начальное значение счетчика NN.
  2. Потоки, вызывающие метод await(), блокируются до тех пор, пока счетчик не станет равен нулю.
  3. Любой поток может вызвать метод countDown(), который уменьшает счетчик на единицу.

CountDownLatch — это одноразовый механизм. Как только счетчик достигает нуля, защелка открывается навсегда. Сбросить счетчик обратно невозможно.

Вернемся к задаче инициализации кэшей. Создадим защелку со счетчиком 3. Главный поток вызовет await(), а рабочие потоки будут вызывать countDown() по мере завершения своей части работы.

public class ServerStartup {
    public static void main(String[] args) throws InterruptedException {
        CountDownLatch latch = new CountDownLatch(3);

        startCacheLoader("DB Cache", latch);
        startCacheLoader("Config Cache", latch);
        startCacheLoader("Redis Cache", latch);

        System.out.println("Главный поток ожидает загрузки кэшей...");
        // Блокируемся, пока счетчик не станет 0
        latch.await();
        System.out.println("Все кэши загружены. Сервер готов к приему запросов.");
    }

    private static void startCacheLoader(String name, CountDownLatch latch) {
        new Thread(() -> {
            try {
                System.out.println(name + " загружается...");
                Thread.sleep((long) (Math.random() * 1000)); // Имитация работы
                System.out.println(name + " загружен.");
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            } finally {
                // Гарантированно уменьшаем счетчик, даже если была ошибка
                latch.countDown();
            }
        }).start();
    }
}

Обратите внимание на блок finally. Это критически важное правило: вызов countDown() должен происходить при любом исходе операции. Если один из потоков упадет с исключением и не уменьшит счетчик, главный поток останется заблокированным в await() навечно.

У CountDownLatch есть и обратный паттерн использования — «стартовый пистолет». Защелка инициализируется значением N=1N = 1. Множество рабочих потоков вызывают await() и ждут. Главный поток выполняет подготовку, а затем вызывает countDown(), одновременно освобождая все рабочие потоки. Это полезно для нагрузочного тестирования, когда нужно создать максимальный всплеск конкурентности.

CyclicBarrier: точка сбора

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

Здесь CountDownLatch не подойдет, так как его нельзя переиспользовать. Для таких сценариев создан CyclicBarrier (циклический барьер).

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

  1. При создании барьера указывается количество потоков (parties), которые должны к нему подойти.
  2. Поток, достигнув контрольной точки, вызывает метод await() и блокируется.
  3. Когда последний поток вызывает await(), барьер «прорывается»: все ожидающие потоки разблокируются, а сам барьер автоматически сбрасывается в исходное состояние для следующего цикла.

Дополнительно конструктор CyclicBarrier принимает объект Runnable — действие, которое будет выполнено ровно один раз последним пришедшим потоком перед тем, как все потоки будут разблокированы. Это идеальное место для объединения результатов (merge).

public class VideoProcessor {
    public static void main(String[] args) {
        int threadsCount = 4;

        // Барьер на 4 потока + действие при достижении барьера
        CyclicBarrier barrier = new CyclicBarrier(threadsCount, () -> {
            System.out.println("--- Кадр полностью обработан. Склеиваем результаты. ---");
        });

        for (int i = 0; i < threadsCount; i++) {
            int segmentId = i;
            new Thread(() -> processFrames(segmentId, barrier)).start();
        }
    }

    private static void processFrames(int segmentId, CyclicBarrier barrier) {
        try {
            for (int frame = 1; frame <= 3; frame++) {
                System.out.println("Поток " + segmentId + " обрабатывает сегмент кадра " + frame);
                Thread.sleep((long) (Math.random() * 1000));

                System.out.println("Поток " + segmentId + " ждет остальных на барьере...");
                // Поток блокируется, пока не придут остальные 3
                barrier.await();
            }
        } catch (Exception e) {
            Thread.currentThread().interrupt();
        }
    }
}

Если один из потоков, ожидающих на CyclicBarrier, прерывается (вызывается interrupt()) или падает по таймауту, барьер считается сломанным (broken). Все остальные потоки, ожидавшие на этом барьере, немедленно получат BrokenBarrierException. Это защищает систему от взаимной блокировки, если один из участников «не дошел до встречи».

Главное концептуальное отличие

Начинающие разработчики часто путают эти два синхронизатора. Чтобы навсегда закрыть этот вопрос, запомните одно правило:

CountDownLatch ждет событий (вызовов countDown). CyclicBarrier ждет потоков (вызовов await).

В CountDownLatch один поток может вызвать countDown() несколько раз, уменьшая счетчик. Защелке неважно, кто уменьшает счетчик, ей важно, сколько раз это произошло. В CyclicBarrier счетчик привязан к физическим потокам. Один поток не может пройти барьер за двоих.

Характеристика CountDownLatch CyclicBarrier
Что считает События (операции) Потоки
Переиспользование Одноразовый Циклический (сбрасывается автоматически)
Действие при срабатывании Нет Можно передать Runnable для выполнения
Кто блокируется Тот, кто вызвал await() (обычно 1 главный поток) Все участники (рабочие потоки) ждут друг друга

Оба этих синхронизатора отлично справляются со своими задачами, когда мы работаем с сырыми потоками. Но в современных Java-приложениях мы редко создаем потоки вручную через new Thread(). Обычно мы отправляем задачи в пулы потоков и хотим получать результаты их выполнения. Для этого требуются другие инструменты, к которым мы и перейдем на следующем этапе.

Почему нельзя вручную создавать потоки на каждый запрос

Почему нельзя вручную создавать потоки на каждый запрос

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

while (true) {
    Socket socket = serverSocket.accept();
    new Thread(() -> processRequest(socket)).start();
}

Эта архитектура называется «один поток на запрос» (Thread-per-Request). Она отлично работает, когда у вас 10, 100 или даже 500 одновременных пользователей. Код выглядит элегантно, потоки изолированы, синхронизация минимальна. Но как только сервер сталкивается с серьезной нагрузкой — скажем, 10 000 одновременных запросов — приложение неизбежно «ложится».

Чтобы понять, почему мы вынуждены отказаться от ручного управления потоками и перейти к пулам (пакет java.util.concurrent), разберем три фундаментальных барьера, о которые разбивается паттерн new Thread().

1. Накладные расходы на создание (Latency)

Как мы выяснили в самых первых главах, объект Thread в Java — это лишь тонкая обертка над реальным потоком операционной системы (Platform Thread).

Вызов start() не просто создает Java-объект в куче. Он инициирует системный вызов (JNI), заставляя ядро ОС выделить структуры данных, настроить регистры процессора и выделить память под стек. На современных системах создание и уничтожение одного потока занимает около 1–2 миллисекунд.

Если ваша задача processRequest — это быстрая операция (например, отдать статический файл из кэша за 1 миллисекунду), то получается абсурдная ситуация: вы тратите больше времени на создание и уничтожение потока, чем на выполнение самой полезной работы. Вычислительные ресурсы тратятся на обслуживание инфраструктуры, а не на бизнес-логику.

2. Предел по памяти (Out Of Memory)

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

Каждый поток должен иметь собственное изолированное пространство для хранения локальных переменных и фреймов вызовов методов — стек (Stack). По умолчанию в 64-битных JVM размер стека для одного потока составляет 1 МБ (регулируется флагом -Xss).

Эта память выделяется вне Java Heap (кучи), в так называемой Native Memory (нативной памяти ОС). Зависимость потребления памяти от количества потоков линейна: Memory=N×StackSizeMemory = N \times StackSize

Если на сервер приходит 5000 одновременных запросов, ОС должна выделить 5 ГБ оперативной памяти только под пустые стеки потоков, не считая памяти под сами объекты приложения. Когда свободная оперативная память заканчивается, JVM выбрасывает фатальную ошибку, которую невозможно обработать:

java.lang.OutOfMemoryError: unable to create new native thread

3. Предел по процессору: Переключение контекста

Допустим, мы арендовали сервер с 256 ГБ оперативной памяти. Теперь память не проблема. Сможем ли мы эффективно запустить 100 000 потоков на 16-ядерном процессоре? Нет.

Процессор с 16 ядрами физически может выполнять только 16 потоков одновременно. Чтобы создать иллюзию параллельной работы 100 000 потоков, планировщик ОС (Scheduler) вынужден постоянно приостанавливать одни потоки и запускать другие. Этот процесс называется переключением контекста (Context Switch).

При переключении контекста процессор должен:

  1. Сохранить текущие значения всех регистров в память.
  2. Очистить конвейер инструкций и сбросить кэши процессора (L1/L2), так как данные старого потока новому не нужны.
  3. Загрузить в регистры состояние нового потока из памяти.

Переключение контекста стоит дорого (тысячи тактов процессора). Когда активных потоков становится слишком много, возникает эффект пробуксовки (Thrashing): операционная система начинает тратить почти 100% процессорного времени на переключение между потоками, а на выполнение полезного кода времени не остается вообще.

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

Смена парадигмы: от потоков к задачам

Проблема паттерна new Thread(() -> task).start() заключается в жестком связывании двух совершенно разных понятий:

  1. Задача (Task) — логическая единица работы (обработать HTTP-запрос, посчитать хэш).
  2. Исполнитель (Thread) — физический ресурс системы (поток ОС).

Создавая поток на каждый запрос, мы позволяем внешнему миру (количеству пользователей) напрямую диктовать, сколько физических ресурсов должна выделить наша система. Это прямой путь к отказу в обслуживании (DoS).

Чтобы построить отказоустойчивую и масштабируемую систему, нам необходимо разорвать эту связь. Задачи должны складываться в очередь, а обрабатывать их должна строго ограниченная группа переиспользуемых потоков. Именно эту фундаментальную задачу решает архитектура пулов потоков, которую мы разберем в следующей главе.

Архитектура ExecutorService и стандартные пулы потоков

Архитектура ExecutorService и стандартные пулы потоков

В предыдущей главе мы выяснили, что создание нового потока на каждую задачу — прямой путь к падению приложения. Если на сервер обрушится 10 000 запросов в секунду, попытка создать 10 000 потоков приведет либо к исчерпанию памяти (OOM), либо к параличу процессора из-за постоянного переключения контекста.

Нам нужен механизм, который позволит отобразить MM поступающих задач на NN постоянно работающих потоков, где NMN \ll M. В Java эту проблему решает фреймворк Executor и концепция пула потоков.

Разделение задачи и исполнителя

До появления пулов потоков разработчики сами управляли жизненным циклом Thread. Код был жестко связан: задача (Runnable) и механизм ее выполнения (Thread) создавались одновременно.

Фреймворк, появившийся в java.util.concurrent, ввел фундаментальное архитектурное изменение: отделение задачи от механизма ее выполнения.

В основе лежат три ключевых интерфейса:

  1. Executor — базовый интерфейс с единственным методом execute(Runnable command). Он ничего не знает о жизненном цикле потоков. Его единственная гарантия: вы передаете задачу, и она будет выполнена (возможно, в другом потоке, а возможно, и в вызывающем).
  2. ExecutorService — расширяет Executor, добавляя методы для управления жизненным циклом самого пула (shutdown(), shutdownNow()) и методы для отправки задач с возможностью отслеживания результата (submit()).
  3. ThreadPoolExecutor — конкретная реализация ExecutorService, которая инкапсулирует в себе всю сложную логику управления пулом потоков, очередями и лимитами.

Переход на ExecutorService означает, что вы перестаете писать new Thread(task).start(). Вместо этого вы пишете pool.execute(task), делегируя JVM решение о том, какой именно поток возьмет эту задачу и когда.

Анатомия ThreadPoolExecutor: как это работает под капотом

Чтобы писать предсказуемый код, нужно понимать, как именно ThreadPoolExecutor принимает решение: создать новый поток, положить задачу в очередь или отклонить ее.

У пула есть четыре главных параметра конфигурации:

  • corePoolSize — базовое количество потоков. Эти потоки создаются по мере поступления задач и не уничтожаются (по умолчанию), даже если простаивают.
  • maximumPoolSize — максимальное количество потоков. Пул может расшириться до этого значения при пиковых нагрузках.
  • workQueue — блокирующая очередь, в которой задачи ждут, пока освободится рабочий поток.
  • keepAliveTime — время жизни «лишних» потоков (тех, что сверх corePoolSize). Если такой поток простаивает дольше этого времени, он завершается, освобождая ресурсы ОС.

Неочевидный алгоритм добавления задачи

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

Когда вы вызываете execute(task), пул действует по строгому алгоритму:

  1. Проверка Core: Если текущее количество потоков меньше corePoolSize, пул создает новый поток для этой задачи (даже если другие потоки прямо сейчас свободны).
  2. Очередь: Если количество потоков достигло corePoolSize, пул пытается положить задачу в workQueue.
  3. Расширение до Max: Если очередь заполнена, и задача туда не помещается, пул создает новый поток, пока их общее количество не достигнет maximumPoolSize.
  4. Отказ: Если достигнут maximumPoolSize, а очередь всё ещё полна, пул вызывает RejectedExecutionHandler (по умолчанию выбрасывается исключение RejectedExecutionException).

Этот алгоритм делает очередь первичным буфером, а расширение пула до максимума — крайней мерой перед отказом в обслуживании.

Фабрика Executors и стандартные пулы

Чтобы не настраивать ThreadPoolExecutor вручную каждый раз, в Java есть утилитный класс Executors с фабричными методами. Рассмотрим три самых популярных.

Метод фабрики Core Max Очередь Применение
newFixedThreadPool(n) nn nn LinkedBlockingQueue (безлимитная) Стабильная нагрузка, известное потребление ресурсов.
newCachedThreadPool() 00 \infty SynchronousQueue (вместимость 0) Множество короткоживущих задач, неравномерная нагрузка.
newSingleThreadExecutor() 11 11 LinkedBlockingQueue (безлимитная) Строго последовательное выполнение задач, гарантия порядка.

Почему Executors часто запрещают в production-коде

Несмотря на удобство, использование методов класса Executors во многих компаниях считается антипаттерном. Причина кроется в их скрытых настройках ThreadPoolExecutor.

Проблема FixedThreadPool:

public static ExecutorService newFixedThreadPool(int nThreads) {
    return new ThreadPoolExecutor(nThreads, nThreads,
                                  0L, TimeUnit.MILLISECONDS,
                                  new LinkedBlockingQueue<Runnable>());
}

Он использует LinkedBlockingQueue без указания размера. Это значит, что ее вместимость равна Integer.MAX_VALUE. Если задачи поступают быстрее, чем nn потоков успевают их обрабатывать (например, зависла база данных), очередь будет расти бесконечно. Это приведет к исчерпанию памяти (Out Of Memory Error).

Проблема CachedThreadPool:

public static ExecutorService newCachedThreadPool() {
    return new ThreadPoolExecutor(0, Integer.MAX_VALUE,
                                  60L, TimeUnit.SECONDS,
                                  new SynchronousQueue<Runnable>());
}

Здесь очередь SynchronousQueue вообще не хранит элементы (вместимость 0), но зато maximumPoolSize равен бесконечности. Если придет 10 000 задач одновременно, пул создаст 10 000 потоков. Мы возвращаемся к проблеме переключения контекста и падения JVM, от которой пытались уйти.

Решение: В высоконагруженных системах всегда создавайте ThreadPoolExecutor вручную. Это заставляет вас явно ответить на вопросы: сколько потоков я могу себе позволить? Какого размера очередь мне нужна? Что делать, если очередь переполнится?

// Пример безопасного пула для production
ExecutorService safePool = new ThreadPoolExecutor(
    10, // corePoolSize
    50, // maximumPoolSize
    60L, TimeUnit.SECONDS, // keepAliveTime
    new ArrayBlockingQueue<>(1000) // Ограниченная очередь!
);

Мы научились безопасно передавать задачи на выполнение пулу потоков через метод execute(Runnable). Но интерфейс Runnable имеет ограничение: его метод run() возвращает void и не может выбросить проверяемое исключение. Если нам нужно поручить пулу вычислить сложный отчет и получить этот отчет обратно, Runnable не поможет. Для этого потребуются новые абстракции, которые мы разберем далее.

Получение результатов вычислений с помощью Callable и Future

Получение результатов вычислений с помощью Callable и Future

В прошлой главе мы научились эффективно распределять задачи по пулам потоков с помощью ExecutorService. Но до сих пор мы отправляли в пул только объекты Runnable. Метод run() этого интерфейса возвращает void.

Представьте, что вы отправили в пул задачу: рассчитать криптографический хэш файла, сделать HTTP-запрос к стороннему API или выполнить сложный SQL-запрос. Как главному потоку получить результат этой работы? Использовать общие переменные и синхронизировать доступ к ним через synchronized или wait/notify? Это вернет нас к низкоуровневой рутине, от которой мы пытались уйти, используя пулы.

Для элегантного решения этой проблемы в Java существуют интерфейс Callable и абстракция Future.

Callable: задача, которая возвращает результат

Интерфейс Callable<V> был добавлен в Java 5 специально для работы с пулами потоков. В отличие от Runnable, он параметризован типом возвращаемого значения V и его единственный метод call() имеет право выбрасывать проверяемые исключения (checked exceptions).

Сравним два интерфейса:

Характеристика Runnable Callable<V>
Метод void run() V call() throws Exception
Возвращаемый результат Нет (void) Есть (объект типа V)
Обработка исключений Не может выбрасывать checked exceptions Может выбрасывать любые исключения
Использование Фоновые задачи (очистка, логирование) Вычисления, запросы к БД, сетевые вызовы

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

import java.util.concurrent.Callable;

public class FetchUserProfileTask implements Callable<String> {
    private final int userId;

    public FetchUserProfileTask(int userId) {
        this.userId = userId;
    }

    @Override
    public String call() throws Exception {
        // Симуляция задержки сети
        Thread.sleep(2000);
        return "User_Profile_" + userId;
    }
}

Чтобы отправить эту задачу в пул, мы используем метод submit() вместо execute(). Но метод submit() не блокирует текущий поток, он мгновенно возвращает управление. Как же получить строку "User_Profile_...", если она будет готова только через две секунды?

Future: долговая расписка на данные

Когда вы передаете Callable в метод submit(), пул потоков возвращает вам объект типа Future<V>.

Future (будущее) — это своеобразная «долговая расписка». Представьте, что вы заказали бургер в ресторане быстрого питания. Кассир принял оплату и выдал вам чек с номером заказа. Ваш бургер еще не готов, но у вас на руках есть объект (чек), с помощью которого вы сможете забрать еду, когда повар закончит работу.

В коде это выглядит так:

ExecutorService executor = Executors.newFixedThreadPool(2);

// Отправляем задачу, мгновенно получаем "расписку"
Future<String> future = executor.submit(new FetchUserProfileTask(42));

System.out.println("Задача отправлена, делаем другие дела...");

// ... здесь главный поток может выполнять другую полезную работу ...

У интерфейса Future есть несколько ключевых методов для управления результатом:

  • isDone() — возвращает true, если задача завершена (успешно, с ошибкой или отменена). Не блокирует поток.
  • get() — возвращает результат вычисления. Блокирует вызывающий поток, если результат еще не готов.
  • get(long timeout, TimeUnit unit) — ждет результат ограниченное время, выбрасывает TimeoutException, если время вышло.
  • cancel(boolean mayInterruptIfRunning) — пытается отменить выполнение задачи.

Блокирующая природа метода get()

Самый важный нюанс работы с Future кроется в методе get(). Когда главный поток вызывает future.get(), он проверяет статус задачи. Если задача еще выполняется рабочим потоком в пуле, главной поток переходит в состояние ожидания (WAITING) и засыпает, пока результат не будет готов.

Именно поэтому вызов get() нужно делать только тогда, когда результат действительно необходим для продолжения работы, предварительно выполнив все возможные параллельные действия.

Чтобы избежать вечной блокировки (например, если сетевой запрос внутри Callable завис), в production-коде всегда следует использовать версию get с таймаутом:

try {
    // Ждем максимум 3 секунды
    String result = future.get(3, TimeUnit.SECONDS);
    System.out.println("Получен результат: " + result);
} catch (TimeoutException e) {
    System.err.println("Задача не успела выполниться за 3 секунды!");
    future.cancel(true); // Отменяем зависшую задачу
}

Обработка исключений: куда исчезают ошибки?

Вспомним проблему из первых глав: если поток падает с исключением внутри run(), исключение выводится в консоль (через UncaughtExceptionHandler), а поток умирает. Главный поток об этом даже не узнает.

С Callable и Future механика кардинально меняется. Если внутри метода call() выбрасывается исключение (любое — RuntimeException или проверяемое), пул потоков перехватывает его. Рабочий поток не умирает. Пул сохраняет это исключение внутри объекта Future.

Когда вы вызовете future.get(), это сохраненное исключение будет обернуто в ExecutionException и выброшено в вызывающий поток.

Callable<Integer> failingTask = () -> {
    return 10 / 0; // ArithmeticException
};

Future<Integer> future = executor.submit(failingTask);

try {
    future.get(); // Здесь вылетит ExecutionException
} catch (ExecutionException e) {
    // Чтобы получить реальную причину ошибки, используем getCause()
    Throwable rootCause = e.getCause();
    System.out.println("Задача упала из-за: " + rootCause.getClass().getSimpleName());
    // Выведет: Задача упала из-за: ArithmeticException
} catch (InterruptedException e) {
    // Выбрасывается, если ГЛАВНЫЙ поток был прерван во время ожидания get()
    Thread.currentThread().interrupt();
}

Отмена задачи: как работает cancel()

Метод future.cancel(boolean mayInterruptIfRunning) позволяет прервать задачу. Поведение зависит от статуса задачи и переданного флага:

  1. Задача еще в очереди (не начала выполняться): Она просто удаляется из очереди и никогда не запустится. Флаг mayInterruptIfRunning не имеет значения.
  2. Задача уже выполняется, cancel(false): Пул позволит задаче завершиться естественным путем, но вызовы get() у этого Future сразу выбросят CancellationException.
  3. Задача уже выполняется, cancel(true): Пул вызовет метод interrupt() у рабочего потока, который в данный момент выполняет эту задачу.

Здесь мы возвращаемся к концепции кооперативной отмены, которую разбирали во второй главе. Вызов cancel(true) не убивает поток принудительно. Он лишь выставляет флаг прерывания. Чтобы задача действительно остановилась, код внутри call() должен реагировать на прерывание: либо вызывать блокирующие методы вроде Thread.sleep() (которые выбросят InterruptedException), либо регулярно проверять Thread.currentThread().isInterrupted().

Пакетная обработка: метод invokeAll

Часто возникает потребность выполнить не одну задачу, а целый пакет, и дождаться результатов их всех. Например, для отрисовки главной страницы интернет-магазина нужно параллельно сходить в три микросервиса: получить профиль пользователя, его корзину и рекомендации.

Вместо того чтобы вручную вызывать submit в цикле и собирать Future в список, интерфейс ExecutorService предоставляет метод invokeAll.

Метод invokeAll принимает коллекцию Callable и блокирует вызывающий поток, пока все задачи не завершатся (успешно или с ошибкой). На выходе он возвращает список Future, в котором гарантированно для каждого элемента метод isDone() вернет true.

List<Callable<String>> tasks = Arrays.asList(
    () -> fetchFromService("Auth"),
    () -> fetchFromService("Cart"),
    () -> fetchFromService("Recommendations")
);

try {
    // Главный поток блокируется здесь, пока все 3 задачи не завершатся
    List<Future<String>> futures = executor.invokeAll(tasks);

    for (Future<String> f : futures) {
        // Вызов get() здесь отработает мгновенно, без блокировки,
        // так как invokeAll уже дождался завершения всех задач
        System.out.println("Результат: " + f.get());
    }
} catch (InterruptedException e) {
    Thread.currentThread().interrupt();
}

Callable и Future решили проблему получения результатов из пула потоков и элегантной обработки ошибок. Однако у Future есть серьезный архитектурный недостаток: мы не можем сказать «когда результат будет готов, передай его в следующую функцию». Нам приходится либо активно опрашивать isDone(), либо блокировать поток вызовом get(). Эта проблема решается построением асинхронных конвейеров, о которых мы поговорим в следующей главе.

Асинхронное программирование с CompletableFuture

Асинхронное программирование с CompletableFuture

В прошлой главе мы научились отправлять задачи в пул потоков и получать объект Future. Но у классического Future есть фатальный архитектурный изъян: метод get() блокирует вызывающий поток. Представьте высоконагруженный сервер, где пул из 200 потоков обрабатывает входящие запросы. Если каждый поток отправит запрос к базе данных и вызовет Future.get(), все 200 потоков уснут в ожидании ответа. Пул исчерпан, сервер перестает принимать новые соединения, хотя процессор при этом абсолютно свободен.

Чтобы утилизировать ресурсы эффективно, нам нужен механизм, который говорит: «Не жди результата. Вот тебе инструкция (коллбек) — когда данные придут, выполни её». В Java эту парадигму реализует CompletableFuture.

От контейнера к конвейеру: интерфейс CompletionStage

CompletableFuture реализует не только знакомый нам Future, но и интерфейс CompletionStage (стадия завершения). Это меняет всё. Теперь мы работаем с асинхронной операцией не как с черным ящиком, из которого нужно силой вытягивать результат, а как со звеном в цепочке обработки данных.

Создадим простую задачу: загрузить строку из сети. Вместо ручного создания Callable и отправки его в пул, мы используем фабричный метод:

CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
    return downloadData(); // Долгая сетевая операция
});

Пока данные скачиваются в фоновом потоке, вызывающий поток (например, main) мгновенно переходит к следующей строке кода. Теперь мы можем прикрепить к этому future инструкции о том, что делать дальше:

future.thenApply(data -> data.toUpperCase())     // Трансформация (как map)
      .thenAccept(result -> System.out.println(result)); // Потребление (как forEach)

Мы построили конвейер. Методы thenApply и thenAccept не блокируют текущий поток. Они лишь регистрируют функции, которые будут вызваны когда-нибудь потом, как только downloadData() завершит работу.

Кто выполняет коллбеки? Правило суффикса Async

Ключевой вопрос для написания эффективного многопоточного кода: в каком именно потоке выполнится лямбда, переданная в thenApply?

По умолчанию CompletableFuture.supplyAsync() использует общий пул потоков JVM — ForkJoinPool.commonPool() (его устройство мы детально разберем в будущих главах). Но что происходит с методами-коллбеками?

Здесь действует важное правило: коллбек без суффикса Async выполняется в том потоке, который завершил предыдущую стадию, либо в вызывающем потоке, если стадия уже завершена.

Рассмотрим пример:

CompletableFuture.supplyAsync(() -> fetchUser(), dbPool)
                 .thenApply(user -> enrichWithAnalytics(user));

Если fetchUser() выполняется долго, то поток из dbPool, завершив загрузку, сам пойдет выполнять enrichWithAnalytics(). Но если fetchUser() отработал мгновенно (например, данные были в кэше) еще до того, как главный поток успел вызвать thenApply, то enrichWithAnalytics() будет выполнен прямо в главном потоке (main)!

Если enrichWithAnalytics() — тяжелая вычислительная задача, выполнение её в пуле базы данных (или в main) может нарушить архитектуру и привести к нехватке потоков БД. Чтобы жестко задать, где должен выполняться следующий шаг, используйте методы с суффиксом Async:

CompletableFuture.supplyAsync(() -> fetchUser(), dbPool)
                 .thenApplyAsync(user -> enrichWithAnalytics(user), computePool);

Теперь, независимо от скорости выполнения первой стадии, вторая стадия гарантированно будет отправлена в виде новой задачи в computePool.

Композиция задач: FlatMap для многопоточности

Самая мощная сторона CompletableFuture раскрывается, когда нужно скоординировать несколько асинхронных операций.

Представьте сценарий: мы асинхронно получаем ID пользователя, а затем по этому ID нужно асинхронно запросить его корзину покупок. Если мы используем thenApply, произойдет неприятность:

// Плохо: вложенные CompletableFuture
CompletableFuture<CompletableFuture<Cart>> nested =
    CompletableFuture.supplyAsync(() -> getUserId())
                     .thenApply(id -> fetchCartAsync(id));

Метод fetchCartAsync сам возвращает CompletableFuture. В итоге мы получаем "матрешку", с которой крайне неудобно работать.

Для цепочек, где каждый следующий шаг сам является асинхронным и зависит от результата предыдущего, используется thenCompose (аналог flatMap из Stream API):

// Хорошо: плоская структура
CompletableFuture<Cart> flat =
    CompletableFuture.supplyAsync(() -> getUserId())
                     .thenCompose(id -> fetchCartAsync(id));

thenCompose дождется завершения внутреннего CompletableFuture и передаст его результат дальше по конвейеру, сохраняя структуру плоской.

Объединение независимых задач

Другой частый сценарий: нужно сделать два независимых запроса параллельно, дождаться обоих и объединить результаты. Например, запросить цену товара в одном сервисе, а персональную скидку — в другом.

Для этого используется метод thenCombine:

CompletableFuture<Double> priceFuture = CompletableFuture.supplyAsync(() -> getPrice());
CompletableFuture<Double> discountFuture = CompletableFuture.supplyAsync(() -> getDiscount());

CompletableFuture<Double> finalPrice = priceFuture.thenCombine(discountFuture, (price, discount) -> {
    return price - (price * discount);
});

Здесь priceFuture и discountFuture выполняются параллельно. Лямбда внутри thenCombine будет вызвана только тогда, когда оба фьючерса завершатся успешно.

Обработка исключений в конвейере

В классическом Future исключения инкапсулировались и выбрасывались как ExecutionException при вызове get(). В асинхронном конвейере get() не вызывается, поэтому ошибки текут по конвейеру вниз, пропуская все успешные блоки (подобно тому, как исключение летит по стеку вызовов до ближайшего catch).

Для перехвата ошибок есть метод exceptionally. Он позволяет вернуть резервное значение, если что-то пошло не так:

CompletableFuture.supplyAsync(() -> fetchUnreliableData())
    .thenApply(data -> process(data))
    .exceptionally(ex -> {
        System.err.println("Ошибка: " + ex.getMessage());
        return "Дефолтные данные"; // Восстановление после ошибки
    })
    .thenAccept(result -> System.out.println(result));

Если fetchUnreliableData или process выбросят исключение, конвейер перепрыгнет сразу к exceptionally. Если всё пройдет без ошибок, exceptionally будет просто проигнорирован.

Если же вам нужно обработать и успешный результат, и ошибку в одном месте (например, чтобы гарантированно закрыть ресурсы или залогировать статус), используется метод handle:

.handle((result, ex) -> {
    if (ex != null) {
        log.error("Провал", ex);
        return "Fallback";
    }
    return result;
})

CompletableFuture превращает многопоточный код из набора блокировок в декларативный граф вычислений. Мы описываем зависимости между данными, а JVM сама решает, когда и в каком потоке эти данные обрабатывать, оставляя ресурсы свободными для других задач.

Потокобезопасные коллекции: ConcurrentHashMap и CopyOnWriteArrayList

Потокобезопасные коллекции: ConcurrentHashMap и CopyOnWriteArrayList

Представьте высоконагруженный кэш, из которого сотни потоков ежесекундно читают настройки системы, и лишь изредка один поток их обновляет. Если мы используем обычный HashMap и обернем его в Collections.synchronizedMap(), мы получим потокобезопасность, но убьем производительность. Синхронизирующая обертка использует одну глобальную блокировку на весь объект: пока один поток читает значение, остальные 99 потоков (даже те, что хотят прочитать другие ключи) выстраиваются в очередь и переходят в состояние BLOCKED.

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

CopyOnWriteArrayList: читаем без блокировок

Когда речь заходит о списках, часто возникает паттерн «часто читаем, редко пишем». Типичный пример — список слушателей событий (listeners) или кэш конфигураций.

Для таких сценариев в Java есть CopyOnWriteArrayList. Его главная идея заложена в названии: копирование при записи.

Внутри этого класса данные хранятся в обычном массиве, ссылка на который помечена ключевым словом volatile (оно гарантирует, что изменения ссылки мгновенно станут видны всем потокам).

Алгоритм работы выглядит так:

  1. Чтение (get, iterator) происходит абсолютно без блокировок. Поток просто берет текущую ссылку на внутренний массив и читает из него данные. Сложность операции — O(1)O(1).
  2. Запись (add, set, remove) использует блокировку ReentrantLock. Поток-писатель захватывает замок, создает новую копию внутреннего массива (размером на один элемент больше), добавляет элемент в эту копию, а затем атомарно перезаписывает volatile ссылку, подменяя старый массив новым. После этого замок освобождается.

Что происходит с потоками, которые начали итерироваться по списку прямо перед тем, как писатель подменил массив? Ничего страшного. Они продолжат читать старую копию массива, на которую у них осталась локальная ссылка. Итератор CopyOnWriteArrayList никогда не выбрасывает ConcurrentModificationException, потому что он работает со слепком данных на момент своего создания.

Цена такого подхода: Каждая операция изменения требует выделения памяти и копирования всего массива. Если в списке миллион элементов, добавление одного элемента потребует скопировать миллион ссылок. Сложность записи составляет O(N)O(N). Поэтому CopyOnWriteArrayList категорически не подходит для сценариев, где данные часто меняются или список имеет огромный размер.

ConcurrentHashMap: эволюция блокировок

Если список можно скопировать целиком, то копировать хэш-таблицу при каждом добавлении ключа — непозволительная роскошь. ConcurrentHashMap (CHM) решает задачу иначе: он дробит блокировки.

Как это было раньше (Lock Striping)

В старых версиях Java (до 8-й) CHM использовал технику Lock Striping (полосатая блокировка). Вся хэш-таблица делилась на 16 независимых сегментов, каждый из которых защищался собственным замком. Если два потока писали в разные сегменты, они делали это параллельно. Это было лучше глобальной блокировки, но при высокой конкуренции 16 сегментов все равно становились узким местом.

Современный подход (Java 8+)

Начиная с Java 8, ConcurrentHashMap был полностью переписан. Разработчики отказались от сегментов в пользу блокировки на уровне отдельной «корзины» (bucket) — конкретной ячейки внутреннего массива.

Как теперь происходит операция put(key, value):

  1. Вычисляется хэш ключа и определяется индекс корзины в массиве.
  2. Если корзина пуста, CHM использует операцию CAS (Compare-And-Swap) для установки первого узла. CAS — это оптимистичная атомарная инструкция на уровне процессора, которая позволяет изменить значение переменной без использования тяжелых блокировок ОС. (Детально механику CAS мы разберем в главе про Lock-Free алгоритмы).
  3. Если корзина не пуста (там уже есть элементы, образующие связный список или красно-черное дерево), поток берет synchronized блокировку только на первый узел этой конкретной корзины.

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

Операция get(key) в современном CHM работает полностью без блокировок. Узлы внутри корзин используют volatile поля для ссылок на следующие элементы и сами значения, что гарантирует видимость изменений для читающих потоков сразу после того, как писатель завершил работу.

Ловушка Check-Then-Act в потокобезопасных коллекциях

Использование ConcurrentHashMap гарантирует, что сама структура данных не сломается (не возникнет бесконечных циклов при рехэшировании, не потеряются узлы). Но она не гарантирует правильность вашей бизнес-логики.

Рассмотрим классический пример инициализации кэша. Нам нужно получить значение по ключу, а если его нет — вычислить, положить в мапу и вернуть.

// АНТИПАТТЕРН: Состояние гонки
if (!cache.containsKey(key)) {
    Value val = expensiveComputation();
    cache.put(key, val);
}
return cache.get(key);

Даже если cache — это ConcurrentHashMap, этот код содержит состояние гонки (Race Condition). Два потока могут одновременно проверить containsKey (оба получат false), оба выполнят тяжелое вычисление expensiveComputation(), и затем по очереди запишут результат в мапу. Потокобезопасность коллекции гарантирует лишь то, что метод put отработает корректно, но не может запретить двум потокам зайти внутрь блока if.

Правильное решение: атомарные методы

Чтобы решить эту проблему, ConcurrentHashMap предоставляет методы, которые выполняют сложную логику внутри своей внутренней блокировки корзины.

Вместо связки «проверь, затем положи», нужно использовать метод computeIfAbsent:

// ПРАВИЛЬНО: Атомарная операция
return cache.computeIfAbsent(key, k -> expensiveComputation());

Как это работает под капотом:

  1. CHM находит нужную корзину.
  2. Захватывает монитор (synchronized) на первом узле корзины.
  3. Проверяет, есть ли ключ.
  4. Если ключа нет, прямо под блокировкой выполняет переданную лямбду expensiveComputation().
  5. Записывает результат, отпускает монитор и возвращает значение.

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

Итоги

Потокобезопасные коллекции из пакета java.util.concurrent — это мощный инструмент, который позволяет избежать ручного управления критическими секциями.

  • Используйте CopyOnWriteArrayList, когда данные читаются на порядки чаще, чем изменяются, и объем данных невелик.
  • Используйте ConcurrentHashMap как универсальную замену HashMap в многопоточной среде, но помните: потокобезопасность методов коллекции не делает последовательность вызовов этих методов атомарной. Всегда используйте встроенные методы вроде computeIfAbsent или putIfAbsent для составных операций.

Очереди блокирующего типа (BlockingQueue) в шаблоне Producer-Consumer

Очереди блокирующего типа (BlockingQueue) в шаблоне Producer-Consumer

В главе о методах wait и notify мы вручную создавали ограниченный буфер для обмена сообщениями между потоками. Нам приходилось управлять мониторами, проверять условия в циклах while для защиты от ложных пробуждений и жонглировать флагами. Писать такой код в production-системах — всё равно что собирать двигатель автомобиля с нуля перед каждой поездкой на работу.

В реальном мире для безопасной передачи данных между потоками используется шаблон Producer-Consumer (Производитель-Потребитель), а стандартным инструментом для его реализации в Java выступает интерфейс BlockingQueue.

Зачем нужен Producer-Consumer и что такое Backpressure

Представьте систему обработки заказов. Web-потоки принимают HTTP-запросы от пользователей (Производители) и складывают их в очередь. Фоновые потоки (Потребители) забирают заказы из очереди и сохраняют их в базу данных.

Главная ценность этого шаблона — развязка во времени и скорости (decoupling). Производители и Потребители ничего не знают друг о друге и могут работать с разной скоростью.

Но что произойдет, если в Черную пятницу HTTP-запросы начнут поступать со скоростью 10 000 в секунду, а база данных способна обработать только 1 000? Если очередь безлимитная, мы быстро получим OutOfMemoryError. Нам нужен механизм, который заставит Web-потоки притормозить, пока база данных не разгребет завал.

Этот механизм называется Backpressure (обратное давление). И именно слово Blocking в названии BlockingQueue обеспечивает это давление из коробки: очередь умеет усыплять потоки-производители, если она переполнена, и усыплять потоки-потребители, если она пуста.

Матрица методов BlockingQueue

Интерфейс BlockingQueue предоставляет четыре набора методов для добавления и извлечения элементов. Выбор правильного метода критически важен — он определяет, как поведет себя поток в экстремальной ситуации.

Поведение при переполнении/пустоте Добавление (Insert) Извлечение (Remove) Изучение (Examine)
Выбрасывает исключение add(e) remove() element()
Возвращает спец. значение offer(e) poll() peek()
Блокирует поток put(e) take() нет
Блокирует с таймаутом offer(e, time, unit) poll(time, unit) нет
  1. Группа исключений: Если вызовете add(e) для заполненной очереди, получите IllegalStateException. Используется редко, так как ломает логику плавного управления нагрузкой.
  2. Группа специальных значений: Метод offer(e) вернет false, если места нет. Метод poll() вернет null, если очередь пуста. Отличный выбор, если при переполнении вы хотите сразу отбросить задачу или записать метрику, не задерживая поток.
  3. Группа блокировок: Методы put(e) и take() — сердце паттерна Producer-Consumer. Если очередь полна, put переведет поток в состояние WAITING, пока не появится место. Если пуста, take усыпит потребителя до появления элемента.
  4. Группа таймаутов: Поток блокируется, но если время вышло, а операция не удалась, возвращается false или null. Защищает систему от вечных зависаний.

Анатомия очередей: Array vs Linked

В пакете java.util.concurrent есть две основные реализации, которые на первый взгляд делают одно и то же, но радикально отличаются под капотом.

ArrayBlockingQueue: кольцевой буфер и единый замок

Эта очередь построена на основе обычного массива фиксированного размера. При создании вы обязаны указать ее вместимость (capacity). Массив используется как кольцевой буфер (Ring Buffer): два указателя отслеживают начало и конец данных, двигаясь по кругу.

Главная архитектурная особенность ArrayBlockingQueue — использование одного объекта ReentrantLock для защиты всего состояния.

В любой момент времени с очередью может работать либо только один Производитель, либо только один Потребитель. Одновременное добавление и извлечение невозможно.

Это создает высокую конкуренцию (contention) за блокировку, если потоков много. Однако у нее есть огромный плюс: после создания массива очередь больше не выделяет память под новые объекты (нет аллокаций), что снижает нагрузку на Garbage Collector.

LinkedBlockingQueue: два замка и узлы

Эта очередь построена на основе связного списка. Каждый элемент оборачивается во внутренний объект Node. По умолчанию ее размер равен Integer.MAX_VALUE (почти безлимитная), но всегда следует задавать жесткий лимит при создании, чтобы избежать OOM.

Ее архитектурный козырь — использование двух независимых ReentrantLock: putLock и takeLock.

Производители конкурируют только с производителями, а потребители — только с потребителями. Добавление элемента в конец списка и извлечение из начала могут происходить параллельно!

Минус: на каждый добавленный элемент создается новый объект Node, что дает дополнительную работу сборщику мусора.

Парадокс нулевой вместимости: SynchronousQueue

Существует особая реализация, которая ломает привычное понимание очереди — SynchronousQueue. Ее вместимость математически равна нулю: capacity=0capacity = 0.

В ней нет массива или связного списка для хранения элементов. Это точка рандеву (встречи) для потоков. Когда Производитель вызывает put(e), он блокируется и ждет, пока какой-нибудь Потребитель не вызовет take(). Как только они встречаются, происходит прямая передача данных (Direct Handoff) из стека одного потока в стек другого, минуя кучу очереди.

Именно SynchronousQueue используется под капотом пула Executors.newCachedThreadPool(). Если приходит новая задача, пул пытается отдать ее через offer(). Если свободных потоков-потребителей прямо сейчас нет (никто не ждет в take()), очередь возвращает false, и пул принимает решение создать новый поток.

Практический пример: Сервис отправки Email

Соберем полученные знания в надежный сервис отправки уведомлений. Мы применим LinkedBlockingQueue для высокой пропускной способности и методы put/take для реализации Backpressure.

import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;

public class EmailNotificationService {
    // Ограничиваем очередь 1000 письмами для защиты памяти
    private final BlockingQueue<Email> queue = new LinkedBlockingQueue<>(1000);

    public EmailNotificationService() {
        // Запускаем 3 потока-потребителя
        for (int i = 0; i < 3; i++) {
            Thread consumer = new Thread(this::consumeLoop, "Email-Sender-" + i);
            consumer.setDaemon(true);
            consumer.start();
        }
    }

    // Вызывается Web-потоками (Производители)
    public void sendNotification(Email email) {
        try {
            // Если писем накопилось 1000, Web-поток заблокируется здесь (Backpressure)
            queue.put(email);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new RuntimeException("Отправка прервана", e);
        }
    }

    // Фоновый цикл Потребителей
    private void consumeLoop() {
        try {
            while (!Thread.currentThread().isInterrupted()) {
                // Если очередь пуста, поток засыпает и не тратит CPU
                Email email = queue.take();
                sendViaSmtp(email);
            }
        } catch (InterruptedException e) {
            // Корректное завершение при остановке приложения
            Thread.currentThread().interrupt();
        }
    }

    private void sendViaSmtp(Email email) {
        // Симуляция долгой сетевой отправки
        System.out.println(Thread.currentThread().getName() + " отправил: " + email.subject());
    }
}

record Email(String to, String subject, String body) {}

В этом коде нам не нужно думать о синхронизации, wait или notify. Очередь LinkedBlockingQueue инкапсулирует всю сложность управления состоянием BLOCKED и WAITING, позволяя нам сосредоточиться на бизнес-логике.

В следующих главах мы шагнем за пределы блокировок и узнаем, как работают коллекции, которые вообще не усыпляют потоки, используя Lock-Free алгоритмы и инструкции процессора.

Концепция неблокирующей синхронизации (Lock-Free)

Концепция неблокирующей синхронизации (Lock-Free)

Представьте оживленный перекресток. Традиционная синхронизация (через synchronized или ReentrantLock) — это светофор. Пока горит красный, поток обязан остановиться и ждать, даже если на пересекаемой дороге никого нет. Неблокирующая синхронизация — это перекресток с круговым движением. Машины не останавливаются в ожидании зеленого света; они притормаживают, оценивают обстановку и вливаются в поток. Если кто-то уже занимает нужную полосу, водитель делает еще один круг и пробует снова.

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

Цена пессимизма

Все механизмы, которые мы рассматривали ранее (мониторы объектов, явные блокировки), реализуют пессимистичный подход к конкуренции. Пессимизм заключается в предположении: «Если я не заблокирую данные прямо сейчас, кто-то другой обязательно их испортит».

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

  1. Переключение контекста (Context Switch): Когда поток не может захватить блокировку, операционная система переводит его в состояние ожидания, сохраняя регистры процессора и очищая кэши. Это дорогая операция.
  2. Взаимная блокировка (Deadlock): Ошибка в порядке захвата замков навсегда останавливает часть системы.
  3. Инверсия приоритетов (Priority Inversion): Ситуация, при которой низкоприоритетный поток захватывает блокировку, а затем вытесняется планировщиком ОС. Высокоприоритетный поток, которому нужны эти данные, вынужден простаивать, ожидая потока с низким приоритетом.
  4. Уязвимость к сбоям: Если поток, захвативший Lock, упадет с ошибкой (или зависнет в бесконечном цикле) до вызова unlock(), все остальные потоки, ожидающие этот замок, будут заблокированы вечно.

Смена парадигмы: Оптимистичная конкуренция

Неблокирующие алгоритмы строятся на оптимистичном подходе. Поток исходит из того, что коллизии редки, поэтому он не просит разрешения на работу с данными.

Вместо схемы «Заблокируй \rightarrow Прочитай \rightarrow Измени \rightarrow Разблокируй», оптимистичный поток действует иначе:

  1. Прочитать текущее состояние данных.
  2. Вычислить новое состояние в локальной памяти (в стеке потока).
  3. Попытаться применить изменения.
  4. Если за время вычислений другой поток успел изменить оригинал — отбросить свои вычисления и начать с шага 1.

Главное отличие: при оптимистичном подходе потоки никогда не переводятся ОС в спящий режим. Если возникает коллизия, поток просто повторяет попытку (крутится в цикле), оставаясь активным.

Иерархия гарантий прогресса

В индустрии термин «Lock-Free» часто используют как синоним «код без слова synchronized». Это грубая ошибка. Самописный SpinLock (цикл while(isLocked) {}) не использует системных блокировок, но если поток-владелец флага isLocked зависнет, все остальные потоки зависнут вместе с ним, сжигая процессорное время.

В информатике неблокирующие алгоритмы классифицируются по строгим математическим гарантиям прогресса (Progress Guarantees).

1. Obstruction-Free (Свобода от помех)

Это самая слабая неблокирующая гарантия. Алгоритм является Obstruction-Free, если поток гарантированно завершит свою операцию за конечное число шагов, при условии, что все остальные потоки приостановлены. На практике это означает, что система защищена от зависания потока-владельца (нет дедлоков), но подвержена livelock — ситуации, когда активные потоки постоянно мешают друг другу, заставляя друг друга отбрасывать результаты и начинать заново, из-за чего система в целом стоит на месте.

2. Lock-Free (Свобода от блокировок)

Алгоритм называется Lock-Free, если при одновременном доступе нескольких потоков к разделяемому ресурсу хотя бы один поток гарантированно добивается прогресса (завершает операцию) за конечное число шагов. Это системная гарантия. Система в целом всегда движется вперед. Если 10 потоков столкнулись при попытке обновить значение, 9 из них могут потерпеть неудачу и уйти на повторный круг, но ровно один обязательно добьется успеха. Слабое место: Lock-Free не защищает от голодания (Starvation). Конкретный неудачливый поток может бесконечно проигрывать конкуренцию более быстрым соседям.

3. Wait-Free (Свобода от ожидания)

Святой Грааль многопоточности. Алгоритм является Wait-Free, если каждый поток гарантированно завершает свою операцию за конечное число шагов, независимо от действий и скорости других потоков. Здесь нет голодания по определению. Wait-Free алгоритмы критически важны для систем реального времени (авиация, медицинское оборудование), где время отклика должно быть строго детерминировано. Однако писать такие алгоритмы невероятно сложно, и они часто требуют избыточного копирования памяти.

Анатомия Lock-Free прогресса

Давайте посмотрим, как гарантия Lock-Free выглядит в динамике. Поскольку хотя бы один поток всегда побеждает, система работает как воронка: количество конфликтующих потоков постоянно уменьшается.

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

Что делает оптимизм возможным?

Возникает закономерный вопрос: как реализовать шаг «Попытаться применить изменения»? Если мы сначала проверим, не изменились ли данные (if (current == expected)), а затем запишем новые (current = newValue), мы получим классическое состояние гонки (Check-Then-Act). Нам нужна гарантия, что проверка и запись произойдут как единое, неделимое действие, в которое не сможет вклиниться ни один другой поток.

Мы не можем использовать для этого synchronized, иначе алгоритм перестанет быть Lock-Free. Нам нужна помощь на аппаратном уровне. Современные процессоры предоставляют специальные атомарные инструкции, которые позволяют выполнить эту проверку и запись за один такт работы с памятью. О том, как устроена эта магия на уровне железа и как она представлена в Java, мы поговорим в следующей главе.

Сравнение с обменом (CAS) как основа атомарных операций

Сравнение с обменом (CAS) как основа атомарных операций

В предыдущей главе мы выяснили, что Lock-Free алгоритмы строятся на оптимистичной конкуренции: поток читает данные, вычисляет новый результат и пытается применить его. Но здесь возникает парадокс. Если два потока одновременно попытаются применить вычисленный результат, мы снова получим состояние гонки (Race Condition). Нам нужен механизм, который позволит сказать памяти: «Запиши сюда новое значение, но только в том случае, если с момента моего последнего чтения эти данные никто не изменил».

Этот механизм невозможно реализовать исключительно средствами языка программирования. Он требует поддержки на уровне «железа».

Аппаратный фундамент: инструкция cmpxchg

Основой всей неблокирующей синхронизации в Java (и большинстве других языков) является процессорная инструкция Compare-And-Swap (CAS) — «сравнение с обменом».

На архитектуре x86 эта инструкция называется cmpxchg (в сочетании с префиксом LOCK для многопроцессорных систем). Суть в том, что процессор берет на себя гарантию атомарности сложной логической операции.

Как работает CAS концептуально? Операция принимает три аргумента:

  1. Адрес в памяти (где лежит переменная).
  2. Ожидаемое значение (то, которое мы прочитали ранее).
  3. Новое значение (то, которое мы хотим записать).

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

Compare-And-Swap (CAS) — это аппаратная атомарная инструкция, которая обновляет значение в памяти только при условии, что текущее значение совпадает с ожидаемым.

Для программиста эта аппаратная магия выглядит как один неделимый шаг. Невозможно, чтобы контекст потока переключился ровно между этапом сравнения и этапом записи — процессор этого не допустит.

Мост между Java и процессором: Unsafe и VarHandle

Долгое время в Java не было прямого легального способа вызвать инструкцию CAS. Разработчикам JDK приходилось использовать внутренний класс sun.misc.Unsafe, который через JNI (Java Native Interface) транслировал вызовы напрямую в машинный код.

Начиная с Java 9, появился стандартизированный и безопасный инструмент — java.lang.invoke.VarHandle. Это типизированная ссылка на переменную, которая позволяет выполнять над ней низкоуровневые атомарные операции, включая CAS.

В коде вызов CAS обычно выглядит как метод compareAndSet, который возвращает boolean:

// Псевдокод работы с CAS
boolean success = varHandle.compareAndSet(object, expectedValue, newValue);

Если метод вернул true, значит, наш поток оказался первым, и данные успешно обновлены. Если false — нас кто-то опередил, данные в памяти уже другие, и наша попытка записи отклонена.

Сердце Lock-Free: цикл CAS (Spin-wait)

Сам по себе единичный вызов CAS мало полезен. Если он вернул false, мы не можем просто бросить работу и пойти дальше — мы потеряем обновление данных. Поэтому CAS всегда используется внутри цикла while. Этот паттерн называется Spin-wait (активное ожидание).

Рассмотрим анатомию классического CAS-цикла на примере инкремента счетчика:

public void increment() {
    int current;
    int next;
    do {
        // 1. Читаем текущее (актуальное) значение
        current = volatileVariable;

        // 2. Вычисляем новое значение
        next = current + 1;

        // 3. Пытаемся атомарно заменить current на next.
        // Если кто-то успел изменить volatileVariable, CAS вернет false,
        // цикл пойдет на новую итерацию, и мы прочитаем уже новые данные.
    } while (!compareAndSet(current, next));
}

Разберем механику по шагам:

  1. Поток AA читает значение current=5current = 5.
  2. Поток AA вычисляет next=6next = 6.
  3. В этот момент планировщик ОС приостанавливает поток AA.
  4. Поток BB успевает прочитать 55, вычислить 66 и успешно выполнить CAS. В памяти теперь 66.
  5. Поток AA просыпается и вызывает compareAndSet(5, 6).
  6. Процессор видит, что в памяти 66, а ожидалось 55. CAS возвращает false.
  7. Цикл do-while уходит на следующий круг. Поток AA читает свежее значение (66), вычисляет новое (77) и успешно выполняет compareAndSet(6, 7).

Цена оптимизма

Паттерн Spin-wait избавляет нас от тяжелых системных блокировок (мониторов и переключений контекста). Потоки не засыпают, они постоянно находятся в состоянии RUNNABLE.

Однако у этого подхода есть обратная сторона. При высокой конкуренции (High Contention), когда десятки потоков одновременно пытаются обновить одну и ту же переменную через CAS, успешным окажется только один. Остальные получат false и уйдут на новый круг цикла.

В худшем случае потоки могут сотни раз вхолостую крутить цикл while, сжигая процессорное время (CPU time) на постоянное чтение, вычисление и неудачные попытки CAS. Это явление называется удержанием процессора (CPU spinning).

Именно поэтому Lock-Free алгоритмы на базе CAS невероятно эффективны при низкой и средней конкуренции, но могут проигрывать традиционным блокировкам в сценариях, где множество потоков непрерывно и агрессивно атакуют один и тот же участок памяти.

Чтобы не писать циклы do-while и вызовы VarHandle вручную каждый раз, когда нам нужен потокобезопасный счетчик или флаг, в Java предусмотрен пакет java.util.concurrent.atomic. В следующей главе мы разберем, как устроены AtomicInteger и AtomicReference, скрывающие всю магию CAS под капотом.

Атомарные переменные: AtomicInteger, AtomicReference и их устройство

Атомарные переменные: AtomicInteger, AtomicReference и их устройство

В прошлой главе мы выяснили, что инструкция Compare-And-Swap (CAS) в связке с циклом do-while позволяет потокам договариваться об изменении данных без тяжеловесных блокировок. Однако писать низкоуровневые циклы с использованием Unsafe или VarHandle в бизнес-коде — верный путь к ошибкам.

Java инкапсулирует эту механику в пакете java.util.concurrent.atomic. Классы этого пакета предоставляют готовый набор lock-free инструментов. Сегодня мы заглянем под капот AtomicInteger и AtomicReference, чтобы понять, как они превращают аппаратные инструкции в удобный API, и где находится предел их эффективности.

Устройство AtomicInteger изнутри

Если отбросить служебный код, структура AtomicInteger предельно лаконична. Это класс-обёртка над обычным примитивом int, дополненная механизмами прямой работы с памятью.

В исходном коде JDK (начиная с Java 9) это выглядит примерно так:

public class AtomicInteger extends Number implements java.io.Serializable {
    // Инструмент для выполнения низкоуровневых операций (замена Unsafe)
    private static final VarHandle VALUE;

    static {
        try {
            MethodHandles.Lookup l = MethodHandles.lookup();
            VALUE = l.findVarHandle(AtomicInteger.class, "value", int.class);
        } catch (ReflectiveOperationException e) {
            throw new ExceptionInInitializerError(e);
        }
    }

    // Само значение. Обязательно volatile!
    private volatile int value;

    public final int get() {
        return value;
    }
}

Здесь критически важны два элемента:

  1. Поле value объявлено как volatile. Это гарантирует, что любой поток, вызывающий метод get(), мгновенно увидит самое свежее значение, записанное другим потоком. Без volatile поток мог бы прочитать устаревшее значение из кэша своего ядра, и цикл CAS начался бы с неверной предпосылки.
  2. Статический объект VALUE типа VarHandle. Он хранит точное смещение поля value в памяти относительно начала объекта AtomicInteger. Именно через этот дескриптор вызываются нативные CAS-инструкции процессора.

Анатомия методов-мутаторов

Когда вы вызываете atomic.incrementAndGet(), под капотом запускается классический паттерн Spin-wait (активное ожидание), который мы разбирали ранее.

Метод обращается к утилите, которая выполняет следующий цикл:

public final int incrementAndGet() {
    int current;
    int next;
    do {
        current = this.get();       // 1. Читаем актуальное значение
        next = current + 1;         // 2. Вычисляем новое
    } while (!VALUE.compareAndSet(this, current, next)); // 3. Пытаемся применить
    return next;
}

Функциональное обновление: updateAndGet

Простое прибавление единицы — частая, но тривиальная задача. Что если новое значение зависит от сложной бизнес-логики? В Java 8 в атомарные классы добавили методы, принимающие лямбда-выражения, например updateAndGet(IntUnaryOperator updateFunction).

Это позволяет внедрить любую математику внутрь lock-free цикла:

AtomicInteger maxScore = new AtomicInteger(0);

// Атомарно обновляем рекорд, только если новый результат выше
maxScore.updateAndGet(current -> Math.max(current, newResult));

Важное правило: функция, передаваемая в updateAndGet, должна быть чистой (без побочных эффектов). Поскольку при высокой конкуренции CAS может провалиться несколько раз, лямбда-выражение будет выполнено процессором заново для каждой попытки. Если внутри лямбды вы пишете лог в файл или отправляете сетевой запрос, это действие выполнится непредсказуемое количество раз.

AtomicReference и атомарность сложных состояний

AtomicInteger и AtomicLong отлично подходят для счетчиков. Но что делать, если состояние объекта описывается несколькими полями, которые должны меняться строго одновременно?

Допустим, у нас есть банковский счет. У него есть баланс (число) и статус (активен/заблокирован). Мы не можем использовать два отдельных атомика:

  1. Поток А меняет баланс на 0.
  2. Переключение контекста.
  3. Поток Б видит нулевой баланс, но статус всё ещё "активен", и принимает неверное решение.
  4. Поток А меняет статус на "заблокирован".

Атомарность нарушена. Решение — упаковать все связанные поля в один неизменяемый (immutable) объект и переключать ссылку на него с помощью AtomicReference.

class AccountState {
    final int balance;
    final String status;

    AccountState(int balance, String status) {
        this.balance = balance;
        this.status = status;
    }
}

AtomicReference<AccountState> stateRef = new AtomicReference<>(
    new AccountState(1000, "ACTIVE")
);

Теперь мы можем изменить оба поля за одну атомарную операцию, подменив ссылку на объект целиком:

stateRef.updateAndGet(currentState -> {
    if (currentState.balance < amount) {
        throw new InsufficientFundsException();
    }
    int newBalance = currentState.balance - amount;
    String newStatus = (newBalance == 0) ? "BLOCKED" : currentState.status;

    // Создаем НОВЫЙ объект состояния
    return new AccountState(newBalance, newStatus);
});

Инструкция CAS процессора работает только с одним участком памяти размером 64 бита. В случае AtomicReference этим участком является ссылка на объект в куче. Процессор атомарно меняет старый адрес памяти на новый.

Тёмная сторона атомиков: цена высокой конкуренции

Атомарные переменные обеспечивают гарантию прогресса системы (Lock-Free), исключая взаимные блокировки (Deadlocks). Однако они не бесплатны.

Проблема кроется в аппаратной архитектуре. Когда поток успешно выполняет CAS-инструкцию, процессор обязан оповестить все остальные ядра о том, что их кэши для этой ячейки памяти больше недействительны (инвалидация кэша).

Если у нас есть глобальный AtomicInteger, который одновременно пытаются инкрементировать NN потоков на разных ядрах:

  1. Один поток успешно выполняет CAS.
  2. Кэши оставшихся N1N-1 ядер сбрасываются.
  3. N1N-1 потоков вынуждены заново читать значение из основной памяти (или L3 кэша).
  4. Они вычисляют новое значение и снова сталкиваются в CAS-гонке.

При низкой и средней нагрузке AtomicInteger работает в десятки раз быстрее, чем synchronized блок. Но при экстремальной конкуренции (например, метрика, которую обновляют сотни потоков одновременно) графики производительности пересекаются. Потоки тратят почти всё процессорное время на «холостое вращение» в цикле do-while и пересылку сообщений об инвалидации кэшей между ядрами.

Атомарные переменные — идеальный инструмент для управления состояниями (как в примере с AtomicReference) или счетчиками с умеренной конкуренцией. Но для высоконагруженных глобальных счетчиков в Java предусмотрен другой, масштабируемый механизм, который обходит узкое место единой точки синхронизации. О нём мы поговорим в следующей главе.

Масштабируемые счетчики: LongAdder и LongAccumulator

Масштабируемые счетчики: LongAdder и LongAccumulator

В высоконагруженной системе подсчет событий — например, количества HTTP-запросов (RPS) или метрик бизнес-логики — кажется тривиальной задачей. Мы уже знаем, что AtomicLong решает проблему потокобезопасности с помощью аппаратной инструкции CAS. Но при экстремальной конкуренции, когда тысячи потоков одновременно пытаются обновить одну переменную, возникает эффект «бутылочного горлышка». Из NN потоков только один успешно выполняет CAS, а остальные N1N-1 уходят на следующий круг Spin-wait, сжигая процессорное время и вызывая непрерывную инвалидацию кэш-линий.

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

Принцип шардирования: от одной переменной к массиву

Идея, лежащая в основе масштабируемых счетчиков, позаимствована из архитектуры распределенных баз данных — это шардирование (разделение) данных.

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

Именно этот паттерн реализует класс LongAdder, появившийся в Java 8.

Внутри LongAdder устроен умнее, чем просто фиксированный массив. Он адаптируется под текущий уровень конкуренции, чтобы не расходовать память впустую:

  1. Состояние покоя (Base). Изначально LongAdder хранит только одно значение — переменную base. Пока конкуренции нет, потоки успешно обновляют base через обычный CAS, и счетчик работает точно так же, как AtomicLong.
  2. Обнаружение конкуренции (Cells). Как только CAS на переменной base завершается неудачей (что означает, что другой поток нас опередил), LongAdder понимает: началась конкуренция. В этот момент он создает массив объектов Cell (ячеек).
  3. Маршрутизация. Поток, потерпевший неудачу на base, вычисляет хэш и привязывается к конкретному Cell в массиве, выполняя CAS уже над ним.

Каждый объект Cell — это, по сути, урезанная версия AtomicLong, хранящая свое локальное значение. Чтобы потоки, пишущие в соседние ячейки массива, не мешали друг другу на уровне процессорных кэшей, класс Cell снабжен специальной аннотацией @Contended. Она добавляет пустые байты (padding), раздвигая ячейки в памяти так, чтобы они гарантированно попадали в разные кэш-линии.

Чтение результата: компромисс согласованности

Распределенная запись решает проблему производительности. Но как узнать итоговое значение счетчика?

Метод sum() в LongAdder работает по простому алгоритму: он берет значение base и прибавляет к нему значения из всех существующих объектов Cell.

Итог=base+i=0N1CelliИтог = base + \sum_{i=0}^{N-1} Cell_i

Где NN — размер массива ячеек.

Здесь кроется главный архитектурный компромисс LongAdder. Метод sum() не блокирует ячейки во время подсчета. Пока цикл бежит по массиву и суммирует значения, другие потоки могут продолжать инкрементировать как уже пройденные ячейки, так и те, до которых цикл еще не дошел.

Из-за этого LongAdder обеспечивает лишь согласованность в конечном счете (Eventual Consistency). Если вам нужно строгое совпадение счетчика с реальным количеством обработанных элементов для принятия критических бизнес-решений (например, генерация уникальных ID для транзакций), LongAdder не подойдет. Его стихия — сбор метрик, статистика и мониторинг, где погрешность в доли процента в момент чтения не играет роли.

LongAccumulator: обобщение логики

LongAdder умеет только складывать. Но что, если при высокой конкуренции нам нужно находить максимальное значение среди тысяч потоков? Или перемножать числа?

Для произвольных операций существует класс LongAccumulator. Он использует ту же самую механику (переменная base + массив Cell), но вместо жестко зашитого сложения принимает в конструктор функцию — LongBinaryOperator — и начальное значение (identity).

Пример создания аккумулятора, который конкурентно находит максимальное значение:

// Функция Math::max и начальное значение Long.MIN_VALUE
LongAccumulator maxAccumulator = new LongAccumulator(Math::max, Long.MIN_VALUE);

// В разных потоках:
maxAccumulator.accumulate(42);
maxAccumulator.accumulate(15);
maxAccumulator.accumulate(99);

// Получение результата
long currentMax = maxAccumulator.get(); // Вернет 99

При чтении результата метод get() возьмет начальное значение, применит функцию к base, а затем последовательно применит её ко всем ячейкам массива.

Важное математическое ограничение: передаваемая функция обязана быть коммутативной (AB=BAA \circ B = B \circ A) и ассоциативной ((AB)C=A(BC)(A \circ B) \circ C = A \circ (B \circ C)).

Поскольку потоки распределяются по ячейкам хаотично, порядок, в котором значения будут накапливаться в Cell, а затем собираться в методе get(), абсолютно непредсказуем. Сложение и поиск максимума удовлетворяют этим свойствам. А вот вычитание — нет (так как 10551010 - 5 \neq 5 - 10), поэтому использовать вычитание в LongAccumulator нельзя — результат будет зависеть от случайного порядка работы потоков планировщиком ОС.

Правила выбора: AtomicLong против LongAdder

Несмотря на превосходство LongAdder под нагрузкой, он не заменяет AtomicLong полностью. Выбор инструмента зависит от профиля нагрузки вашего приложения.

Характеристика AtomicLong LongAdder
Конкуренция Низкая или средняя Экстремально высокая
Потребление памяти Минимальное (O(1)O(1)) Растет при конкуренции (O(N)O(N) ячеек)
Чтение результата Строго консистентное Eventual Consistency (может отставать)
Операции Сложение, инкремент, CAS-циклы Только сложение (или кастомная функция без CAS)

Если ваш счетчик читается чаще, чем обновляется (например, флаг конфигурации или редкая метрика), AtomicLong будет эффективнее. Если же счетчик обновляется тысячами потоков каждую миллисекунду, а читается раз в 10 секунд потоком мониторинга — LongAdder станет спасением для процессора.

Проблема ABA и её решение с помощью AtomicStampedReference

Проблема ABA и её решение с помощью AtomicStampedReference

Инструкция CAS (Compare-And-Swap) кажется безупречным инструментом для оптимистичной синхронизации. Мы проверяем: если текущее значение равно ожидаемому, значит, никто его не менял, и можно безопасно записать новое. Но равенство значений не гарантирует неизменность истории. Представьте, что вы сдали чемодан в багаж (состояние А). Грузчик вскрыл его, забрал ценные вещи (состояние B), закрыл и поставил на место (состояние А). Вы получаете чемодан, визуально он тот же самый (CAS успешен), но его внутреннее состояние безвозвратно испорчено. В многопоточном программировании эта уязвимость называется проблемой ABA.

Анатомия проблемы ABA

Проблема ABA возникает, когда поток читает значение из разделяемой памяти, затем приостанавливается, а в это время другие потоки изменяют это значение на другое, а затем возвращают обратно к исходному.

Классический сценарий выглядит так:

  1. Поток 1 читает из памяти значение AA.
  2. Поток 1 вытесняется планировщиком ОС.
  3. Поток 2 записывает значение BB.
  4. Поток 2 (или Поток 3) записывает обратно значение AA.
  5. Поток 1 просыпается, выполняет CAS-операцию, сравнивая текущее значение с ожидаемым AA.
  6. CAS возвращает true, так как значения совпадают. Поток 1 продолжает работу, будучи уверенным, что данные не менялись.

Если мы используем AtomicInteger для подсчета лайков, проблема ABA нас не волнует. Если счетчик был 100, кто-то убрал лайк (99), а потом другой добавил (100), для нашей бизнес-логики 100 — это просто число 100. Операция инкремента отработает корректно.

Но если мы работаем с AtomicReference и строим неблокирующие структуры данных, где значения — это ссылки на узлы, хранящие связи с другими узлами, проблема ABA приводит к катастрофе.

Катастрофа в Lock-Free стеке

Рассмотрим классический пример уязвимости — неблокирующий стек (Treiber Stack). В нем каждый узел хранит данные и ссылку на следующий элемент. Вершина стека (head) управляется через AtomicReference.

Пусть в стеке лежат три элемента: Вершина -> Узел_A -> Узел_B -> Узел_C.

Поток 1 хочет извлечь (pop) элемент со стека. Алгоритм извлечения выглядит так:

  1. Прочитать текущую вершину (head = Узел_A).
  2. Прочитать следующий элемент (next = Узел_B).
  3. Выполнить CAS(head, Узел_A, Узел_B).

Поток 1 выполнил первые два шага, запомнил, что ожидаемая вершина — это Узел_A, а новым head должен стать Узел_B. В этот момент поток засыпает.

В дело вступает Поток 2 и совершает серию операций:

  1. Делает pop(): извлекает Узел_A. Стек становится: Узел_B -> Узел_C.
  2. Делает pop(): извлекает Узел_B. Стек становится: Узел_C. (Узел B теперь может быть удален сборщиком мусора или переиспользован).
  3. Делает push(Узел_A): возвращает Узел_A обратно, но теперь он указывает на текущую вершину. Стек становится: Узел_A -> Узел_C.

Просыпается Поток 1. Он выполняет свой отложенный CAS(head, Узел_A, Узел_B). CAS проверяет: равна ли текущая вершина ожидаемому Узел_A? Да, равна! Поток 2 заботливо вернул Узел_A на место. CAS успешно отрабатывает и делает вершиной стека Узел_B.

Но Узел_B был удален из стека! Теперь вершина стека указывает на оторванный от структуры узел, а реальные данные (Узел_C) потеряны. Структура данных полностью разрушена.

Решение — версионирование ссылок

Чтобы решить проблему ABA, недостаточно проверять только саму ссылку. Нужно проверять еще и факт того, что эта ссылка не переназначалась. Для этого к ссылке добавляется «штамп» (stamp) — целочисленный счетчик версии, который увеличивается при каждом изменении.

Теперь CAS должен атомарно проверить два условия:

  1. Текущая ссылка равна ожидаемой ссылке.
  2. Текущий штамп равен ожидаемому штампу.

Если Поток 2 извлечет Узел_A и положит его обратно, ссылка будет той же (Узел_A), но штамп изменится (например, с 1 на 3). Когда Поток 1 попытается выполнить CAS со старым штампом 1, операция закономерно провалится, и поток уйдет на следующий круг цикла do-while, прочитав актуальное состояние стека.

В Java этот механизм реализован в классе AtomicStampedReference.

Устройство и API AtomicStampedReference

AtomicStampedReference не может просто передать процессору два независимых значения (ссылку и int) для выполнения одной инструкции CAS. Архитектура x86 поддерживает инструкцию cmpxchg16b, которая может атомарно сравнивать и менять 16 байт (две 64-битные ссылки), но Java абстрагирует это иначе.

Внутри AtomicStampedReference хранится неизменяемый статический класс Pair, который объединяет ссылку и штамп в один объект.

Когда вы обновляете значение, AtomicStampedReference под капотом создает новый объект Pair и выполняет стандартный CAS уже над ссылкой на этот Pair.

Главный метод для обновления состояния: compareAndSet(V expectedReference, V newReference, int expectedStamp, int newStamp)

Посмотрим, как выглядит безопасный метод pop() для нашего неблокирующего стека с использованием AtomicStampedReference:

public class ABASafeStack<T> {
    private final AtomicStampedReference<Node<T>> head =
        new AtomicStampedReference<>(null, 0);

    public T pop() {
        // Массив из одного элемента для извлечения текущего штампа
        int[] stampHolder = new int[1];

        while (true) {
            // Одновременно получаем и ссылку, и её штамп
            Node<T> currentHead = head.get(stampHolder);
            int currentStamp = stampHolder[0];

            if (currentHead == null) {
                return null; // Стек пуст
            }

            Node<T> nextNode = currentHead.next;

            // CAS проверяет и ссылку, и штамп
            if (head.compareAndSet(currentHead, nextNode,
                                   currentStamp, currentStamp + 1)) {
                return currentHead.value;
            }
            // Если CAS не удался (изменилась ссылка ИЛИ штамп), повторяем цикл
        }
    }
}

Обратите внимание на специфичный API метода get(int[] stampHolder). Поскольку метод в Java не может вернуть два значения одновременно (и ссылку, и штамп), разработчики JDK применили трюк: метод возвращает ссылку как результат, а штамп записывает в массив из одного элемента, который вы передаете в качестве аргумента. Это позволяет прочитать оба значения из одного объекта Pair за одну операцию без риска, что они рассинхронизируются.

AtomicMarkableReference как альтернатива

В пакете java.util.concurrent.atomic есть родственный класс — AtomicMarkableReference. Его механика абсолютно идентична, но вместо целочисленного штампа int он использует логический флаг boolean.

Он не решает проблему ABA в полной мере (так как флаг может смениться с true на false и обратно на true), но он незаменим в алгоритмах, где требуется концепция «логического удаления».

Например, в неблокирующих связных списках физическое удаление узла (переписывание ссылок) — сложная многошаговая операция. Вместо этого узел сначала «маркируется» как удаленный (флаг меняется на true с помощью AtomicMarkableReference), а физическое удаление происходит позже, асинхронно, когда другие потоки наткнутся на маркированный узел при обходе списка.

Разработка собственного lock-free алгоритма на практике

Разработка собственного lock-free алгоритма на практике

Главное ограничение инструкции CAS (Compare-And-Swap), с которым мы работали до сих пор, заключается в том, что она атомарно обновляет только одну ячейку памяти. Но что делать, если логика структуры данных требует одновременного изменения двух ссылок?

В стеке Трайбера (Treiber Stack) нам хватало одного CAS, потому что стек меняется только с одного конца — вершины (head). Но если мы захотим реализовать классическую очередь (FIFO), нам придется работать с двумя указателями: head (откуда забираем) и tail (куда добавляем).

Вставка элемента в конец связного списка состоит из двух шагов:

  1. Привязать новый узел к текущему последнему: tail.next = newNode.
  2. Сдвинуть указатель хвоста на новый узел: tail = newNode.

Сделать эти два шага одной атомарной операцией невозможно. Если поток выполнит первый шаг и будет вытеснен планировщиком ОС, очередь окажется в промежуточном, неконсистентном состоянии. Как другие потоки должны реагировать на это, не прибегая к блокировкам? Ответ на этот вопрос — суть проектирования сложных неблокирующих алгоритмов.

Паттерн «Взаимопомощь» (Helping)

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

Этот паттерн называется взаимопомощью (Helping). Алгоритм проектируется как конечный автомат: любое промежуточное состояние должно быть легко распознаваемым, и любой поток должен иметь достаточно информации, чтобы перевести систему в следующее консистентное состояние.

Мы разберем этот подход на примере алгоритма очереди Майкла-Скотта (Michael-Scott Queue) — именно он лежит в основе стандартной ConcurrentLinkedQueue в Java.

Структура очереди и фиктивный узел (Dummy Node)

Чтобы избежать множества краевых случаев при работе с пустой очередью (когда head и tail указывают на null), в Lock-Free алгоритмах часто используют паттерн Dummy Node.

При создании очереди в нее сразу помещается пустой (фиктивный) узел. Указатели head и tail изначально ссылаются на него. Таким образом, в очереди всегда есть хотя бы один узел, и нам не нужно писать сложную логику инициализации при первой вставке через CAS.

Определим базовую структуру нашего класса:

import java.util.concurrent.atomic.AtomicReference;

public class LockFreeQueue<T> {

    private static class Node<T> {
        final T value;
        final AtomicReference<Node<T>> next;

        Node(T value, Node<T> next) {
            this.value = value;
            this.next = new AtomicReference<>(next);
        }
    }

    private final AtomicReference<Node<T>> head;
    private final AtomicReference<Node<T>> tail;

    public LockFreeQueue() {
        // Создаем Dummy-узел
        Node<T> dummy = new Node<>(null, null);
        this.head = new AtomicReference<>(dummy);
        this.tail = new AtomicReference<>(dummy);
    }
}

Примечание: В отличие от стека Трайбера, здесь мы не используем AtomicStampedReference. В Java при создании new Node<T> выделяется уникальный адрес в куче. Поскольку узлы очереди после извлечения собираются сборщиком мусора и их адреса не переиспользуются немедленно (как это бывает в C++ с ручным управлением памятью), классическая проблема ABA при вставке и удалении здесь не возникает.

Реализация метода enqueue (Вставка)

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

public void enqueue(T value) {
    Node<T> newNode = new Node<>(value, null);

    while (true) {
        Node<T> currentTail = tail.get();
        Node<T> tailNext = currentTail.next.get();

        // Шаг 1: Проверка консистентности прочитанных данных
        if (currentTail == tail.get()) {

            // Шаг 2: Проверяем, не отстает ли указатель tail
            if (tailNext != null) {
                // Очередь в промежуточном состоянии! Другой поток прикрепил узел,
                // но не успел сдвинуть tail. ПОМОГАЕМ ЕМУ:
                tail.compareAndSet(currentTail, tailNext);
            } else {
                // Хвост действительно последний. Пытаемся прикрепить наш узел.
                if (currentTail.next.compareAndSet(null, newNode)) {
                    // Успех! Наш узел прикреплен.
                    // Теперь пытаемся сдвинуть tail на наш узел.
                    tail.compareAndSet(currentTail, newNode);
                    return; // Выходим из цикла
                }
            }
        }
    }
}

Давайте разберем критические моменты этого кода:

  1. Чтение снимка состояния: Мы читаем tail и его next. Сразу после этого мы проверяем currentTail == tail.get(). Если за время чтения next указатель tail уже изменился, наши данные устарели — уходим на новую итерацию (Spin-wait).
  2. Распознавание промежуточного состояния: Если tailNext != null, это значит, что tail указывает не на последний элемент. Какой-то поток уже выполнил compareAndSet(null, newNode), но заснул перед сдвигом tail.
  3. Helping в действии: Увидев tailNext != null, наш поток выполняет tail.compareAndSet(currentTail, tailNext). Он делает работу за уснувший поток! Только после этого он пойдет на следующий виток цикла, чтобы попытаться вставить уже свой элемент.
  4. Безопасность финального шага: Обратите внимание, что второй CAS (tail.compareAndSet(currentTail, newNode)) может завершиться неудачей, если другой поток успел нам «помочь». И это абсолютно нормально — цель достигнута, tail сдвинут, мы можем спокойно выходить из метода.

Реализация метода dequeue (Извлечение)

Извлечение элемента также требует осторожности. Из-за механизма взаимопомощи указатель tail может отставать от реальности, а из-за наличия Dummy-узла реальные данные всегда лежат в узле, следующем за head.

public T dequeue() {
    while (true) {
        Node<T> currentHead = head.get();
        Node<T> currentTail = tail.get();
        Node<T> headNext = currentHead.next.get();

        // Проверка консистентности снимка
        if (currentHead == head.get()) {

            // Проверка на пустоту очереди (или отставание tail)
            if (currentHead == currentTail) {
                if (headNext == null) {
                    return null; // Очередь действительно пуста
                }
                // Очередь не пуста, но tail отстал (указывает на head). Помогаем!
                tail.compareAndSet(currentTail, headNext);
            } else {
                // Очередь не пуста, извлекаем данные
                T value = headNext.value;

                // Пытаемся сдвинуть head на следующий узел
                if (head.compareAndSet(currentHead, headNext)) {
                    // Старый currentHead становится мусором для GC,
                    // а headNext теперь становится новым Dummy-узлом!
                    return value;
                }
            }
        }
    }
}

Здесь раскрывается элегантность паттерна Dummy Node. Когда мы извлекаем элемент, мы не удаляем узел физически. Мы просто берем значение из headNext, а затем сдвигаем head на этот headNext. Бывший узел с данными логически становится новым пустым Dummy-узлом. Предыдущий Dummy-узел теряет все ссылки на себя и убирается сборщиком мусора.

Также обратите внимание на блок взаимопомощи внутри dequeue: если head == tail, но headNext != null, это значит, что другой поток начал вставку, прикрепил узел, но еще не сдвинул tail. Прежде чем пытаться что-то извлечь, мы обязаны помочь ему сдвинуть tail, иначе нарушим логику структуры.

Итоги

Разработка Lock-Free алгоритмов кардинально отличается от программирования с блокировками. Вы не можете просто «защитить» участок кода. Вы должны:

  1. Разбить операцию на минимальные шаги, каждый из которых фиксируется одним CAS.
  2. Спроектировать структуру так, чтобы промежуточные состояния были безопасными и видимыми для всех.
  3. Заложить в логику каждого потока обязанность проверять состояние структуры и помогать завершать чужие зависшие операции.

Именно этот подход делает Lock-Free алгоритмы устойчивыми к зависанию отдельных потоков и гарантирует системный прогресс (тот самый Lock-Free guarantee, который мы обсуждали ранее).

Однако, несмотря на отсутствие явных блокировок, мы все еще опираемся на гарантии видимости изменений между потоками. Почему изменения, сделанные через CAS в одном потоке, мгновенно видны другому? За это отвечает Модель памяти Java (JMM), к изучению которой мы и переходим.

Архитектура процессоров, кэши и когерентность

Архитектура процессоров, кэши и когерентность

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

Но если спуститься на уровень «железа», это предположение разбивается о суровую физику. Современный процессор выполняет инструкцию за доли наносекунды, а поход в оперативную память (RAM) занимает около 100 наносекунд. Если бы ядра процессора каждый раз обращались к RAM напрямую, они бы простаивали 99% времени.

Чтобы процессор не «голодал», между ним и RAM выстроена сложная иерархия памяти. И именно она является физической первопричиной всех странностей многопоточного программирования, с которыми мы столкнемся дальше.

Иерархия памяти и кэш-линии

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

  1. Регистры — память внутри самого вычислительного устройства ядра. Доступ мгновенный (0 циклов).
  2. L1-кэш — приватный кэш ядра, разделенный на кэш инструкций и кэш данных. Доступ занимает 3–4 такта.
  3. L2-кэш — приватный кэш ядра (в некоторых архитектурах разделяется между парой ядер), больше по объему, но медленнее (10–12 тактов).
  4. L3-кэш — общий кэш для всех ядер процессора (30–70 тактов).
  5. RAM (Оперативная память) — общая для всей системы (сотни тактов).

Важнейший нюанс: процессор никогда не читает из оперативной памяти отдельные байты. Данные перемещаются между RAM и кэшами блоками фиксированного размера — кэш-линиями (Cache Lines). В современных процессорах x86 и ARM размер кэш-линии обычно составляет 64 байта.

Кэш-линия — это минимальный квант данных, которым оперирует подсистема памяти процессора. Если вы читаете переменную типа int (4 байта), процессор загрузит в кэш L1 всю 64-байтную линию, в которой лежит этот int, прихватив соседние переменные.

Проблема когерентности кэшей

Наличие приватных кэшей L1 и L2 у каждого ядра создает фундаментальную проблему.

Представьте, что Ядро 1 и Ядро 2 одновременно читают переменную XX (исходно X=0X = 0). Оба ядра загружают кэш-линию с XX в свои L1-кэши. Затем Ядро 1 выполняет инкремент: X=1X = 1. Оно записывает новое значение в свой L1-кэш.

Что в этот момент видит Ядро 2? В его L1-кэше лежит старая копия, где X=0X = 0. Возникает рассинхронизация — нарушение когерентности. Если ничего не предпринять, многопоточное программирование будет невозможным, так как потоки на разных ядрах будут жить в параллельных реальностях.

Протокол MESI: аппаратное решение

Чтобы кэши ядер договаривались между собой, инженеры придумали протоколы когерентности. Самый известный из них — MESI (и его современные вариации вроде MOESI).

Протокол MESI добавляет к каждой кэш-линии два бита состояния. Название — это аббревиатура этих четырех состояний:

  • M (Modified) — линия изменена только в кэше этого ядра. В оперативной памяти лежат устаревшие данные. Ядро имеет эксклюзивное право на запись.
  • E (Exclusive) — линия есть только в кэше этого ядра, и она совпадает с данными в RAM. Ядро может менять её без опроса других ядер.
  • S (Shared) — линия присутствует в кэшах нескольких ядер. Данные совпадают с RAM. Разрешено только чтение.
  • I (Invalid) — данные в этой линии устарели (недействительны). Читать из неё нельзя.

Как это работает в динамике? Ядра общаются друг с другом через специальную внутреннюю шину, рассылая сообщения.

Если Ядро 1 хочет записать данные в линию, которая находится в состоянии Shared, оно обязано разослать по шине сообщение Invalidate всем остальным ядрам. Только получив от них ответы Invalidate Acknowledge (подтверждение, что они перевели свои линии в состояние Invalid), Ядро 1 переводит свою линию в Modified и выполняет запись.

Цена синхронности: Store Buffer и Invalidation Queue

Протокол MESI гарантирует строгую консистентность, но у него есть фатальный недостаток: он синхронный.

Когда Ядро 1 отправляет сообщение Invalidate, оно вынуждено остановить выполнение инструкций и ждать ответов Invalidate Acknowledge от других ядер. В масштабах тактовых частот процессора это целая вечность.

Чтобы не простаивать, архитекторы процессоров добавили две аппаратные оптимизации:

  1. Store Buffer (Буфер записи). Находится между ядром и L1-кэшем. Когда ядро хочет записать данные, оно отправляет Invalidate, но не ждет ответа. Оно кладет новое значение в Store Buffer и продолжает выполнять следующие инструкции. Когда ответы придут, данные из буфера асинхронно «стекут» в L1-кэш.
  2. Invalidation Queue (Очередь инвалидации). Находится на входе в кэш другого ядра. Ядро 2, получив Invalidate, не стирает линию в кэше немедленно (это может занять время, если кэш занят). Оно помещает сообщение в очередь, мгновенно отправляет Acknowledge Ядру 1, и обязуется применить инвалидацию чуть позже.

Иллюзия порядка: как железо ломает логику

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

Рассмотрим классический пример. Исходно A=0A = 0 и B=0B = 0. Ядро 1 выполняет:

  1. A=1A = 1
  2. B=1B = 1

Ядро 2 выполняет:

  1. Читает BB
  2. Читает AA

Предположим, переменная BB уже была в кэше Ядра 1 (состояние Exclusive), а переменная AA была у Ядра 2 (состояние Shared).

Что произойдет физически:

  1. Ядро 1 пишет A=1A = 1. Линия в состоянии Shared, поэтому ядро кладет A=1A = 1 в Store Buffer и отправляет Invalidate Ядру 2.
  2. Ядро 1 не ждет! Оно переходит к B=1B = 1. Линия Exclusive, запись идет прямо в L1-кэш. Состояние становится Modified.
  3. Ядро 2 читает BB. Возникает промах кэша, оно запрашивает BB по шине. Ядро 1 отдает B=1B = 1.
  4. Ядро 2 читает AA. Сообщение Invalidate от Ядра 1 всё ещё лежит в Invalidation Queue Ядра 2 и не применено. Поэтому Ядро 2 читает старое значение A=0A = 0 из своего L1-кэша.

Итог: Ядро 2 увидело, что B=1B = 1, но A=0A = 0. С точки зрения программиста, процессор переупорядочил инструкции (выполнил запись BB до записи AA). Но процессор ничего не переупорядочивал! Инструкции выполнились по порядку, просто данные AA «застряли» в буферах, а данные BB проскочили в кэш сразу.

Барьеры памяти: мост между железом и кодом

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

Разрешается этот конфликт с помощью барьеров памяти (Memory Barriers / Fences). Это специальные процессорные инструкции, которые принудительно синхронизируют буферы.

  • Барьер на запись (Store Barrier) заставляет ядро дождаться, пока все данные из Store Buffer не будут сброшены в L1-кэш.
  • Барьер на чтение (Load Barrier) заставляет ядро дождаться применения всех сообщений из Invalidation Queue, прежде чем читать данные.

Разработчики на Java не пишут ассемблерные инструкции барьеров вручную. Вместо этого JVM берет на себя роль транслятора. Когда мы используем ключевое слово volatile, блоки synchronized или атомарные CAS-операции, JVM автоматически вставляет нужные барьеры памяти для той архитектуры процессора (x86, ARM), на которой сейчас работает приложение.

Именно поэтому изменения, сделанные в AtomicReference в нашем lock-free алгоритме из прошлой главы, мгновенно и корректно видны другим потокам — под капотом JVM расставила барьеры, заставив процессоры сбросить Store Buffers.

В следующей главе мы поднимемся с аппаратного уровня обратно на уровень Java и разберем, как спецификация Java Memory Model (JMM) абстрагирует всю эту аппаратную сложность в элегантную математическую модель отношений Happens-Before.

Спецификация Java Memory Model и отношение Happens-Before

Спецификация Java Memory Model и отношение Happens-Before

В прошлой главе мы заглянули под капот современных процессоров и увидели пугающую картину. Из-за буферов записи (Store Buffers) и очередей инвалидации (Invalidation Queues) процессоры постоянно переупорядочивают инструкции. Если бы нам приходилось вручную расставлять аппаратные барьеры памяти под каждую архитектуру (x86, ARM, PowerPC), принцип Java «написал один раз — работает везде» был бы невозможен.

Чтобы избавить программиста от необходимости знать, как устроен L1-кэш на конкретном процессоре, создатели Java ввели абстракцию. Эта абстракция — Java Memory Model (JMM).

JMM — это не структура памяти в куче (heap) и не стек потока. Это строгий формальный контракт между программистом и JVM. Контракт гласит: если программист соблюдает определенные правила синхронизации, JVM берет на себя обязательство расставить все необходимые аппаратные барьеры памяти так, чтобы программа работала предсказуемо на любом железе.

Иллюзия физического времени

Главная ошибка при написании многопоточного кода — опираться на физическое время.

Допустим, поток А выполняет x=1x = 1 ровно в 12:00:00. Поток B читает переменную xx ровно в 12:00:01. Увидит ли поток B единицу? С точки зрения здравого смысла — да, ведь прошла целая секунда. С точки зрения JMM — нет никаких гарантий.

Значение x=1x = 1 может навсегда остаться в Store Buffer ядра, на котором работал поток А, или поток B может прочитать устаревшее значение из своего локального кэша. В многопоточности физическое время не имеет значения. Значение имеет только видимость (visibility).

Чтобы формализовать эту видимость, спецификация JMM (JSR-133) вводит фундаментальное понятие — отношение Happens-Before (выполняется прежде).

Happens-Before — это логическое отношение между двумя действиями AA и BB. Если действие AA happens-before действия BB (записывается как ABA \to B), то JMM гарантирует: результаты действия AA будут видимы действию BB, и для действия BB будет казаться, что AA выполнилось до него.

Если между записью переменной в одном потоке и её чтением в другом нет связи Happens-Before, возникает Data Race (состояние гонки). В этом случае JVM имеет полное право переупорядочивать инструкции как угодно, и результат чтения становится непредсказуемым.

Базовые правила Happens-Before

Спецификация определяет строгий набор правил, которые создают связи Happens-Before. Рассмотрим самые важные из них.

1. Правило порядка выполнения (Program Order Rule)

Внутри одного потока каждое действие AA happens-before каждого действия BB, которое идет за ним в исходном коде. С точки зрения потока всё выполняется последовательно. JVM может переупорядочивать инструкции внутри потока для оптимизации, но только если это не меняет конечный результат для этого же потока (as-if-serial semantics).

2. Правило монитора (Monitor Lock Rule)

Освобождение блокировки (unlock) happens-before захвату (lock) того же самого монитора. Когда поток А выходит из блока synchronized, а поток B затем входит в блок synchronized по тому же объекту, поток B гарантированно увидит всё, что сделал поток А до выхода из блока.

3. Правило старта потока (Thread Start Rule)

Вызов метода Thread.start() happens-before любому действию внутри запущенного потока. Если главный поток инициализировал переменные, а затем запустил фоновый поток, этот фоновый поток гарантированно увидит подготовленные данные.

4. Правило завершения потока (Thread Join Rule)

Любое действие внутри потока happens-before возврату из метода Thread.join() для этого потока. Если главный поток вызвал worker.join(), то после того как метод вернет управление, главный поток гарантированно увидит все изменения, сделанные внутри worker.

5. Правило volatile-переменной (Volatile Variable Rule)

Запись в volatile переменную happens-before любому последующему чтению этой же переменной. Подробную механику этого правила мы разберем в следующей главе.

Магия транзитивности

Самое мощное свойство Happens-Before — транзитивность. Если ABA \to B, и BCB \to C, то спецификация гарантирует, что ACA \to C.

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

Транзитивность порождает интересную технику, которую часто называют Piggybacking (езда на чужой спине). Мы можем использовать синхронизацию одной переменной, чтобы протолкнуть в память изменения других, не синхронизированных переменных.

Рассмотрим пример:

class Config {
    int timeout = 10;           // Обычная переменная
    String url = "http://a.com"; // Обычная переменная
    final Object lock = new Object();

    // Выполняется Потоком А
    void updateConfig() {
        timeout = 50;           // Действие 1
        url = "http://b.com";   // Действие 2
        synchronized(lock) {    // Действие 3 (lock)
            // Ничего не делаем
        }                       // Действие 4 (unlock)
    }

    // Выполняется Потоком B
    void printConfig() {
        synchronized(lock) {    // Действие 5 (lock)
            System.out.println(url + " " + timeout); // Действие 6
        }                       // Действие 7 (unlock)
    }
}

Давайте проследим цепочку Happens-Before, если Поток А выполнил updateConfig, а затем Поток B выполнил printConfig:

  1. По правилу порядка выполнения (Program Order): Запись timeout и url (Действия 1, 2) \to unlock монитора (Действие 4).
  2. По правилу монитора (Monitor Lock): unlock в Потоке А (Действие 4) \to lock в Потоке B (Действие 5).
  3. По правилу порядка выполнения: lock в Потоке B (Действие 5) \to чтение timeout и url (Действие 6).

Благодаря транзитивности: Действия 1 и 2 \to Действие 6. Поток B гарантированно увидит новые значения timeout = 50 и url = "http://b.com", хотя сами эти переменные не являются volatile, и мы не меняли их внутри synchronized блока! Изменения «проехались на спине» у блокировки монитора.

Итог

Java Memory Model — это юридический договор. Вы, как программист, обязуетесь связывать потоки через правила Happens-Before (используя synchronized, Thread.join, volatile или классы из java.util.concurrent, которые построены на этих же правилах).

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

В следующей главе мы детально разберем правило volatile переменной и узнаем, как именно JVM запрещает компилятору и процессору переупорядочивать инструкции вокруг неё.

Ключевое слово volatile: видимость и запрет переупорядочивания

Ключевое слово volatile: видимость и запрет переупорядочивания

Представьте простейшую задачу: один поток выполняет фоновую работу в цикле, а другой должен дать ему сигнал на остановку. Вы заводите обычный boolean flag = false, первый поток крутится в while (!flag), а второй в нужный момент делает flag = true. Вы запускаете код, второй поток меняет флаг, но первый поток продолжает работать бесконечно.

Почему это происходит? В прошлой главе мы выяснили, что без связи Happens-Before (HB) потоки ничего друг другу не должны. JIT-компилятор видит, что внутри цикла переменная flag не меняется, и в целях оптимизации выносит её чтение за пределы цикла (hoisting). На уровне железа первый поток читает флаг из регистра или L1-кэша, даже не подозревая, что второй поток давно обновил значение в своём Store Buffer.

Чтобы заставить эту схему работать, нам нужен инструмент, который создаст мост Happens-Before между потоками без тяжеловесных блокировок synchronized. Этим инструментом является volatile.

Правило Volatile Variable

Спецификация JMM определяет чёткое правило:

Запись в volatile переменную vv происходит до (happens-before) любого последующего чтения из этой же переменной vv.

Добавление модификатора volatile к нашему флагу кардинально меняет поведение JVM. Компилятор теряет право кэшировать эту переменную в регистрах или выносить её чтение из цикла. Каждое обращение к volatile переменной — это честное обращение к разделяемой памяти.

Однако видимость самого флага — это лишь верхушка айсберга. Истинная сила volatile раскрывается в том, как оно влияет на другие переменные.

Эффект односторонней проницаемости (Roach Motel)

До Java 5 спецификация volatile гарантировала только видимость самой переменной. Это делало модификатор почти бесполезным для сложных задач. Начиная с Java 5 (JSR-133), семантика была усилена: теперь volatile работает как барьер переупорядочивания для всех окружающих инструкций.

Это поведение часто называют семантикой «Roach Motel» (ловушка для тараканов: зайти можно, выйти нельзя).

  1. Volatile Write (Запись): Обычные операции чтения и записи, идущие в коде до записи в volatile переменную, не могут быть переупорядочены так, чтобы оказаться после неё.
  2. Volatile Read (Чтение): Обычные операции чтения и записи, идущие в коде после чтения из volatile переменной, не могут быть переупорядочены так, чтобы оказаться до него.

Почему барьер односторонний? Компилятору и процессору разрешено сдвигать обычные операции «внутрь» зоны между volatile чтением и записью — это безопасно. Но им категорически запрещено выносить операции «наружу», пересекая volatile границу.

Транзитивность: техника Piggybacking

Соединим три правила, которые мы уже знаем:

  1. Program Order: Операции в рамках одного потока выполняются последовательно (логически).
  2. Volatile Variable: Запись в volatile vv happens-before чтение из vv.
  3. Transitivity: Если AA happens-before BB, а BB happens-before CC, то AA happens-before CC.

Это комбо рождает мощнейший паттерн Piggybacking (езда на чужой спине) — безопасную передачу обычных данных через одну volatile переменную.

Рассмотрим пример инициализации конфигурации:

// Поток 1 (Писатель)
configData = loadConfig(); // 1. Обычная запись
isReady = true;            // 2. Volatile запись

// Поток 2 (Читатель)
if (isReady) {             // 3. Volatile чтение
    use(configData);       // 4. Обычное чтение
}

Разберём цепочку Happens-Before:

  • Операция 1 HB Операция 2 (по правилу Program Order).
  • Операция 2 HB Операция 3 (по правилу Volatile Variable).
  • Операция 3 HB Операция 4 (по правилу Program Order).

Следовательно, благодаря транзитивности, Операция 1 HB Операция 4. Когда Поток 2 видит isReady == true, JMM гарантирует, что он увидит и полностью инициализированный объект configData, хотя сам configData не является volatile. Барьер при записи isReady не позволил компилятору переставить инициализацию конфигурации ниже установки флага.

Как это работает на уровне железа

Чтобы обеспечить описанные гарантии JMM, виртуальная машина Java расставляет в машинном коде аппаратные барьеры памяти (Memory Barriers), о которых мы говорили в главе про архитектуру процессоров.

Согласно JSR-133 Cookbook, JVM генерирует следующие барьеры:

  • Перед volatile записью вставляется барьер StoreStore. Он гарантирует, что все предыдущие записи потока будут сброшены в кэш до того, как произойдёт запись в volatile переменную.
  • После volatile записи вставляется барьер StoreLoad. Это самый тяжёлый барьер. Он принудительно очищает Store Buffer процессора, заставляя ядро дождаться подтверждения от всех остальных ядер по протоколу MESI. Именно он делает изменения глобально видимыми.
  • После volatile чтения вставляются барьеры LoadLoad и LoadStore. Они заставляют процессор дождаться применения всех сообщений из Invalidation Queue, прежде чем выполнять следующие инструкции чтения или записи.

Таким образом, абстрактное правило Happens-Before материализуется в конкретные команды управления аппаратными буферами.

Чего volatile НЕ может: иллюзия атомарности

Самая частая ошибка при работе с volatile — попытка использовать его для счётчиков.

volatile int count = 0;
// В нескольких потоках:
count++;

Модификатор volatile гарантирует видимость, но не обеспечивает атомарность составных операций. Операция count++ — это не одно действие, а три (Read-Modify-Write):

  1. Прочитать текущее значение count.
  2. Прибавить единицу.
  3. Записать новое значение обратно.

Даже если переменная volatile, два потока могут одновременно прочитать значение 1010, оба локально вычислить 1111, и оба записать 1111. Одно инкрементирование будет потеряно. Барьеры переупорядочивания здесь не спасают, так как нет конфликта порядка инструкций — есть конфликт одновременного доступа к данным (Data Race).

Для атомарных операций Read-Modify-Write необходимо использовать либо блокировки (synchronized), либо классы из пакета java.util.concurrent.atomic (основанные на инструкциях CAS), которые мы подробно разбирали ранее.

Резюме

Ключевое слово volatile — это легковесный механизм синхронизации, который не блокирует потоки (не вызывает Context Switch).

Используйте volatile, когда:

  • Переменная используется как флаг состояния (завершение работы, статус инициализации).
  • Вам нужно безопасно опубликовать набор обычных переменных, используя volatile флаг как триггер (паттерн Piggybacking).

Не используйте volatile, когда:

  • Новое значение переменной зависит от её предыдущего значения (i++).
  • Инвариант системы включает несколько переменных, которые должны обновляться строго одновременно (здесь нужен synchronized или Lock).

Безопасная публикация объектов и неизменяемость (Immutability)

Безопасная публикация объектов и неизменяемость (Immutability)

Представьте ситуацию: один поток создает объект User user = new User(25), а другой поток читает эту ссылку и вызывает user.getAge(). Может ли второй поток получить значение 0, хотя в коде конструктора явно написано this.age = age? Да, может. Это одна из самых коварных ловушек многопоточности — проблема частично инициализированного объекта.

Иллюзия атомарности оператора new

На уровне исходного кода создание объекта выглядит как одна неделимая операция. Однако для JVM и процессора выражение new User(25) разбивается на несколько независимых шагов:

  1. Выделение памяти под объект (Memory Allocation).
  2. Инициализация полей значениями по умолчанию (для чисел это 00, для ссылок — null).
  3. Выполнение тела конструктора (присвоение this.age = 25).
  4. Запись адреса выделенной памяти в ссылку user.

Вспомним влияние Store Buffer и оптимизаций компилятора: система имеет право переупорядочивать инструкции, если это не меняет логику выполнения внутри одного потока. Шаги 3 и 4 независимы друг от друга с точки зрения потока-создателя. Процессор может сначала записать адрес объекта в переменную user (шаг 4), и только потом выполнить присвоение внутри конструктора (шаг 3).

Если в этот микросекундный зазор вмешается второй поток, он увидит не-null ссылку user, перейдет по ней в память и прочитает дефолтное значение поля age, равное 00. Объект опубликован небезопасно.

Семантика final и барьер Freeze

До Java 1.5 (JSR-133) эта проблема делала невозможным создание надежных неизменяемых объектов без явной синхронизации. Разработчикам приходилось использовать synchronized даже при чтении полей, которые никогда не менялись после создания.

Спецификация JMM ввела жесткое правило для ключевого слова final. Теперь final — это не просто маркер для компилятора «запрещено менять», это полноценная инструкция для подсистемы памяти.

Когда конструктор завершает работу, JMM гарантирует выполнение специального действия — Freeze (заморозка). На аппаратном уровне в конце конструктора вставляется барьер памяти (обычно StoreStore), который гарантирует: все записи в final-поля будут принудительно сброшены в основную память до того, как ссылка на сам объект станет доступна другим потокам.

Если поле age объявлено как final int age, переупорядочивание шагов 3 и 4 становится невозможным. Любой поток, который увидит ссылку на объект, гарантированно увидит и правильное значение его final-полей.

Важное уточнение: гарантия транзитивна. Если final-поле ссылается на массив или другой объект (например, final List<String> items), JMM гарантирует видимость не только самой ссылки на список, но и содержимого этого списка на момент завершения конструктора.

Неизменяемость (Immutability) как абсолютная защита

Опираясь на семантику final, мы приходим к мощнейшему архитектурному паттерну многопоточного программирования — строгой неизменяемости.

Объект считается строго неизменяемым, если выполняются три условия:

  1. Состояние объекта не может быть изменено после создания (отсутствуют сеттеры).
  2. Все поля объекта объявлены как final.
  3. Объект был создан правильно (ссылка this не утекла во время конструирования).

Неизменяемые объекты обладают уникальным свойством: они потокобезопасны по своей природе. Им не нужны блокировки, volatile или атомики. Поскольку их состояние фиксируется барьером Freeze и больше никогда не меняется, состояние гонки (Data Race) физически невозможно — гонка требует хотя бы одной операции записи, а здесь возможны только чтения.

Вы можете свободно кэшировать такие объекты, передавать их между сотнями потоков и публиковать через обычные поля данных без синхронизации. Именно поэтому String, Integer и классы из пакета java.time спроектированы как неизменяемые.

Утечка this (This Escape)

Гарантии final работают только при условии, что объект полностью сконструирован до того, как другие потоки получат к нему доступ. Существует критическая ошибка проектирования, которая пробивает барьер Freeze — утечка ссылки this из конструктора.

Рассмотрим классический антипаттерн:

public class EventListener {
    private final int id;
    private final String name;

    public EventListener(EventDispatcher dispatcher, int id, String name) {
        this.id = id;
        // Утечка this! Объект еще не достроен.
        dispatcher.register(this);
        this.name = name;
    }
}

В момент вызова dispatcher.register(this) конструктор еще не завершил работу. Барьер Freeze еще не пройден. Однако ссылка на объект (this) уже передана во внешнюю систему.

Если EventDispatcher работает в другом потоке, он может мгновенно вызвать метод у зарегистрированного слушателя. Этот поток увидит поле id инициализированным, но поле name будет равно null, несмотря на то, что оно final.

Утечка this может быть неявной. Самые частые сценарии:

  • Запуск нового потока прямо внутри конструктора (new Thread(this).start()).
  • Публикация объекта в публичную статическую коллекцию.
  • Создание анонимного внутреннего класса внутри конструктора (анонимные классы неявно хранят ссылку на внешний класс).

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

public class EventListener {
    private final int id;
    private final String name;

    private EventListener(int id, String name) {
        this.id = id;
        this.name = name;
    }

    public static EventListener createAndRegister(EventDispatcher dispatcher, int id, String name) {
        EventListener listener = new EventListener(id, name);
        // Конструктор завершен, барьер Freeze пройден. Теперь можно публиковать.
        dispatcher.register(listener);
        return listener;
    }
}

Оптимизации компилятора и JVM: Escape Analysis и Lock Coarsening

Оптимизации компилятора и JVM: Escape Analysis и Lock Coarsening

Представьте метод, в котором создаётся новый объект, затем поток захватывает его монитор через synchronized, изменяет пару полей и завершает работу. Сколько тактов процессора и байт памяти будет потрачено на аллокацию в куче, барьеры памяти и захват блокировки? Правильный ответ в современной JVM: ноль.

Мы уже разобрали, как опасна утечка ссылки на не до конца инициализированный объект (This Escape). Но JIT-компилятор (Just-In-Time) смотрит на проблему утечек под другим углом. Если он может математически доказать, что ссылка на объект никогда не покидает пределы вызвавшего его потока, он получает право переписать ваш код, полностью удалив из него абстракции объектов и синхронизации.

Escape Analysis: когда объект остаётся невидимкой

Escape Analysis (анализ достижимости или анализ утекания) — это фаза работы C2 JIT-компилятора, на которой он вычисляет область видимости объекта. Компилятор строит граф потока управления и проверяет, не передаётся ли ссылка на объект в другие потоки, не сохраняется ли она в статические поля и не возвращается ли из метода.

Если объект признаётся локальным для потока (No Escape), JIT применяет две мощные оптимизации.

Lock Elision (Удаление блокировок)

Любая синхронизация имеет цену: даже неконкурентный захват монитора требует атомарной инструкции (CAS) в заголовке объекта (Mark Word). Но если объект не покидает поток, состояние гонки (Data Race) физически невозможно — другим потокам не на чем синхронизироваться.

Классический пример — использование устаревшего StringBuffer внутри метода:

public String buildMessage(String name) {
    StringBuffer buffer = new StringBuffer(); // Синхронизированный класс
    buffer.append("Hello, ");
    buffer.append(name);
    return buffer.toString();
}

Метод append внутри StringBuffer помечен как synchronized. Однако buffer создаётся локально и его ссылка никуда не передаётся (наружу уходит только результат toString(), который является уже другой строкой). JIT-компилятор видит это и на лету вырезает все инструкции захвата и освобождения монитора. Код выполняется так же быстро, как если бы вы использовали несинхронизированный StringBuilder.

Scalar Replacement (Скалярная замена)

Удаление блокировок экономит процессорное время, но объект всё ещё нужно выделять в куче (Heap), а затем собирать Garbage Collector-ом. Здесь вступает в игру вторая оптимизация — скалярная замена.

Если объект не «утекает» из метода, JVM может вообще его не создавать. Вместо выделения памяти под заголовок объекта и его поля в куче, JIT «расщепляет» объект на составляющие примитивы (скаляры) и размещает их прямо в локальных переменных потока — на стеке или в регистрах процессора.

class Point {
    int x, y;
    Point(int x, int y) { this.x = x; this.y = y; }
}

public int calculateDistance() {
    Point p = new Point(10, 20); // Объект не покидает метод
    return p.x * p.x + p.y * p.y;
}

Вместо дорогостоящей аллокации в куче, JIT превратит этот код в эквивалент:

public int calculateDistance() {
    int p_x = 10;
    int p_y = 20;
    return p_x * p_x + p_y * p_y;
}

Эта оптимизация радикально снижает нагрузку на сборщик мусора (GC). Миллионы временных объектов-обёрток, итераторов и DTO, создаваемых внутри методов, живут и умирают на стеке, не оставляя следа в куче.

Lock Coarsening: укрупнение блокировок

Escape Analysis работает только для локальных объектов. Но что делает JIT, если объект уже опубликован и доступен другим потокам, а мы обращаемся к нему в цикле?

Каждый вход в блок synchronized требует установки барьеров памяти (Monitor Enter / Monitor Exit) и обновления состояния в Mark Word. Если поток захватывает и отпускает один и тот же монитор тысячу раз подряд, накладные расходы становятся колоссальными.

Для решения этой проблемы JIT применяет Lock Coarsening (укрупнение блокировок). Компилятор ищет соседние блоки синхронизации, использующие один и тот же монитор, и объединяет их в один большой блок.

Рассмотрим пример с потокобезопасной коллекцией Vector:

public void addAll(Vector<String> vector, String[] items) {
    for (int i = 0; i < items.length; i++) {
        vector.add(items[i]); // synchronized метод внутри
    }
}

Без оптимизации поток бы захватывал и отпускал монитор vector на каждой итерации цикла. JIT-компилятор анализирует граф выполнения и трансформирует логику (на уровне машинного кода) примерно так:

public void addAll(Vector<String> vector, String[] items) {
    synchronized(vector) { // Блокировка вынесена за пределы цикла
        for (int i = 0; i < items.length; i++) {
            vector.add_without_lock(items[i]);
        }
    }
}

Укрупнение блокировок работает не только с циклами, но и с последовательными вызовами. Если вы пишете obj.syncMethod1(); obj.syncMethod2();, JIT объединит их под один захват монитора.

Важный нюанс: JIT укрупняет блокировки только если между ними нет длительных операций (например, ввода-вывода или Thread.sleep()). Если внутри цикла происходит обращение к сети, JIT оставит блокировки раздельными, чтобы не лишать другие потоки доступа к монитору на долгое время.

Философия JIT и ручные микрооптимизации

Понимание этих механизмов меняет подход к написанию кода. Многие разработчики пытаются применять ручные микрооптимизации:

  • Пишут собственные пулы объектов (Object Pool) для мелких структур данных, чтобы «избежать нагрузки на GC».
  • Избегают создания короткоживущих объектов внутри методов.
  • Пишут громоздкий код, пытаясь вручную вынести блокировки за пределы циклов там, где это нарушает инкапсуляцию.

В современной Java такие действия часто дают обратный эффект. Пул объектов требует синхронизации (или CAS) для выдачи и возврата объекта, а сами объекты в пуле живут долго и попадают в Old Generation кучи. В то же время простой new Point(), благодаря Escape Analysis и Scalar Replacement, не стоит вообще ничего и не требует работы GC.

Пишите идиоматичный, читаемый код с минимальной областью видимости переменных. Чем меньше область видимости объекта, тем проще JIT-компилятору доказать, что объект не утекает, и применить агрессивные оптимизации, превратив ваш высокоуровневый Java-код в сверхбыстрый машинный монолит.

Однако оптимизации JIT не спасают от проблем на уровне архитектуры процессора. Если потоки интенсивно модифицируют независимые переменные, которые случайно оказались рядом в памяти, возникает эффект ложного разделения. Как размер кэш-линии влияет на многопоточность и зачем нужна аннотация @Contended, мы разберём далее.

Влияние ложного разделения кэш-линий (False Sharing) и аннотация @Contended

Влияние ложного разделения кэш-линий (False Sharing) и аннотация @Contended

Представьте простую задачу: два потока собирают метрики. Первый поток инкрементирует счетчик успешных запросов, второй — счетчик ошибок. Оба счетчика лежат в одном объекте:

class Metrics {
    public volatile long successCount = 0;
    public volatile long errorCount = 0;
}

Потоки абсолютно независимы. Они не используют блокировок, не читают данные друг друга и пишут в разные переменные. Казалось бы, мы достигли идеальной параллельности. Но если запустить этот код под нагрузкой, мы обнаружим парадокс: производительность может упасть в 3–5 раз по сравнению с однопоточным выполнением.

Причина кроется не в коде на Java, а в физической архитектуре процессора. Эта проблема называется ложным разделением (False Sharing).

Пинг-понг на системной шине

В главе об архитектуре процессоров мы выяснили, что минимальный квант обмена данными между оперативной памятью и кэшами L1/L2 — это кэш-линия (Cache Line). На большинстве современных архитектур (x86, ARM) размер кэш-линии составляет 64 байта.

Процессор никогда не загружает из памяти отдельную переменную. Если ядру нужен long (8 байт), оно загрузит блок в 64 байта, в котором лежит этот long, а заодно и соседние данные.

Посмотрим на наш класс Metrics. Две переменные типа long занимают 16 байт. В куче (Heap) JVM они будут расположены в памяти вплотную друг к другу. С вероятностью, близкой к 100%, обе переменные попадут в одну и ту же кэш-линию.

Что происходит, когда два потока, работающие на разных физических ядрах, начинают писать в эти переменные?

  1. Ядро 1 модифицирует successCount. По протоколу когерентности MESI, кэш-линия в L1-кэше Ядра 1 переходит в состояние Modified. Процессор рассылает сигнал инвалидации другим ядрам.
  2. Ядро 2 хочет инкрементировать errorCount. Оно обращается к своему L1-кэшу, но видит, что кэш-линия помечена как Invalid (из-за действия Ядра 1).
  3. Ядро 2 вынуждено ждать, пока актуальная кэш-линия сбросится из Ядра 1 в L3-кэш или оперативную память, чтобы загрузить ее себе. Загрузив линию, Ядро 2 переводит ее в состояние Modified, инвалидируя копию у Ядра 1.
  4. Ядро 1 снова хочет обновить successCount и получает промах кэша (Cache Miss), так как линия инвалидирована Ядром 2.

Этот процесс называется Cache Line Ping-Pong. Ядра непрерывно перехватывают друг у друга одну и ту же кэш-линию, пересылая данные по системной шине.

Разделение называется ложным, потому что потоки логически не разделяют данные (каждый работает со своей переменной), но аппаратно они разделяют один и тот же контейнер — кэш-линию.

Ручное выравнивание (Padding) до Java 8

Как разорвать эту связь? Нужно сделать так, чтобы successCount и errorCount физически оказались в разных кэш-линиях. Если размер линии 64 байта, нам нужно вставить между переменными «пустышку» такого же размера.

До появления Java 8 разработчикам высоконагруженных систем приходилось применять технику Padding (заполнение) вручную:

class Metrics {
    public volatile long successCount = 0;

    // Padding: 7 long-переменных = 56 байт
    // Плюс сама successCount (8 байт) = 64 байта
    private long p1, p2, p3, p4, p5, p6, p7;

    public volatile long errorCount = 0;
}

Добавляя 7 неиспользуемых переменных типа long, мы гарантируем, что расстояние между successCount и errorCount составит минимум 64 байта (с учетом заголовка объекта). Теперь переменные лежат в разных кэш-линиях. Ядро 1 может эксклюзивно владеть линией с successCount, а Ядро 2 — линией с errorCount. Протокол MESI больше не заставляет их конфликтовать.

Однако ручной Padding имел серьезные недостатки:

  1. Зависимость от платформы. Мы жестко закладываемся на размер линии в 64 байта. Если код запустится на архитектуре с линией 128 байт, False Sharing вернется.
  2. Борьба с JIT-компилятором. Умная JVM может заметить, что переменные p1...p7 нигде не читаются и не пишутся. Оптимизатор (Dead Code Elimination) имеет полное право просто удалить их из структуры класса, сломав наше выравнивание. Разработчикам приходилось писать фиктивные методы, читающие эти переменные, чтобы обмануть компилятор.

Аннотация @Contended

Чтобы избавить разработчиков от «хаков» с памятью, в Java 8 был реализован JEP 142, который добавил внутреннюю аннотацию @Contended (в пакете jdk.internal.vm.annotation).

Эта аннотация дает команду JVM: «Данная переменная подвержена высокой конкуренции. Размести ее в памяти так, чтобы она гарантированно занимала отдельную кэш-линию».

class Metrics {
    @Contended
    public volatile long successCount = 0;

    @Contended
    public volatile long errorCount = 0;
}

Когда JVM (при загрузке класса) видит эту аннотацию, она сама опрашивает процессор о размере кэш-линии и автоматически вставляет нужное количество пустых байт (padding) до и после поля. JIT-компилятор знает об этой аннотации и никогда не удалит эти отступы.

Более того, @Contended позволяет группировать поля (Contention Groups). Если несколько переменных всегда обновляются одним потоком одновременно, их выгодно держать в одной кэш-линии:

class PlayerState {
    // Координаты обновляются вместе — кладем в одну группу "location"
    @Contended("location")
    public volatile double x;
    @Contended("location")
    public volatile double y;

    // Здоровье обновляется независимо — кладем в другую группу
    @Contended("health")
    public volatile int hp;
}

Важное ограничение для разработчиков

Поскольку аннотация добавляет невидимые байты, она увеличивает потребление памяти (Memory Footprint). Если бездумно вешать её на все поля, куча быстро переполнится пустотой.

Именно поэтому разработчики JDK сделали аннотацию внутренней и отключили её действие для пользовательского кода по умолчанию. Если вы просто напишете @Contended в своем проекте, она будет проигнорирована.

Чтобы JVM начала обрабатывать эту аннотацию в ваших классах, приложение необходимо запускать со специальным флагом виртуальной машины:

java -XX:-RestrictContended MyApp

(Обратите внимание на минус перед Restrict: мы отключаем ограничение).

Возвращение к LongAdder

В главе о масштабируемых счетчиках мы изучали класс LongAdder, который решает проблему конкуренции за счет массива ячеек (Cell[]). Потоки хэшируются и пишут в разные ячейки массива.

Теперь, понимая механику False Sharing, мы можем увидеть скрытую угрозу в дизайне LongAdder. Объекты в массиве лежат в памяти последовательно. Если класс Cell будет занимать всего 24 байта (заголовок объекта + один long), то в одну 64-байтную кэш-линию поместятся сразу 2-3 ячейки.

Если Поток 1 пишет в Cell[0], а Поток 2 пишет в Cell[1], они попадут в одну кэш-линию. Возникнет False Sharing, и вся идея шардирования счетчика провалится — потоки снова начнут блокировать друг друга на уровне железа.

Именно поэтому, если мы заглянем в исходный код java.util.concurrent.atomic.Striped64 (базовый класс для LongAdder), мы увидим следующее:

@jdk.internal.vm.annotation.Contended
static final class Cell {
    volatile long value;
    // ... методы CAS
}

Разработчики JDK пометили весь класс Cell аннотацией @Contended. Это заставляет JVM раздвигать объекты Cell в памяти так, чтобы каждый из них гарантированно начинался с новой кэш-линии. Таким образом, потоки, пишущие в разные ячейки, работают с физически независимыми участками кэша L1, достигая идеальной линейной масштабируемости.

Когда это действительно нужно применять?

False Sharing — это микрооптимизация. О ней стоит задумываться только если выполняются все три условия:

  1. Переменные обновляются очень часто (миллионы раз в секунду).
  2. Переменные обновляются разными потоками.
  3. Профилировщик (например, JMH с профилировщиком perfnorm) показывает аномально высокое количество кэш-промахов (L1-dcache-load-misses).

В 99% бизнес-кода на Java ложное разделение не является узким местом. Однако понимание этого механизма отличает программиста, который просто использует API, от инженера, который понимает, как его код исполняется на кремнии.

Шаблоны проектирования многопоточных приложений (Thread Pool, Work Stealing)

Шаблоны проектирования многопоточных приложений (Thread Pool, Work Stealing)

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

В этой главе мы поднимемся с уровня процессорных кэшей на уровень архитектурных паттернов. Мы разберем, как правильно структурировать распределение задач, почему классический Thread Pool перестает масштабироваться на десятках ядер, и как шаблон Work Stealing решает эту проблему, элегантно используя особенности аппаратного обеспечения.

Паттерн Thread Pool: больше, чем просто переиспользование потоков

Ранее мы уже работали с ExecutorService как с удобным инструментом. Но с точки зрения архитектуры, паттерн Thread Pool решает фундаментальную задачу: он отделяет прием задачи от ее выполнения (Decoupling).

Вместо паттерна «Thread-Per-Message» (создание нового потока на каждый чих), Thread Pool реализует макро-версию шаблона Producer-Consumer:

  1. Клиентские потоки (Producers) генерируют задачи и складывают их в буфер.
  2. Пул рабочих потоков (Consumers) непрерывно извлекает задачи из буфера и выполняет их.

Это дает предсказуемость: мы ограничиваем потребление памяти (количество потоков фиксировано) и избегаем накладных расходов на системные вызовы создания потоков ОС.

Анатомия узкого горлышка

В классическом Thread Pool (например, ThreadPoolExecutor в Java) буфером выступает единая глобальная очередь, чаще всего LinkedBlockingQueue или ArrayBlockingQueue.

Пока у вас 4 или 8 ядер, эта архитектура работает прекрасно. Но представьте сервер с 64 ядрами, где все 64 рабочих потока пытаются одновременно взять новую задачу из одной глобальной очереди.

Даже если мы заменим блокирующую очередь на сверхбыструю неблокирующую (например, алгоритм Майкла-Скотта, который мы писали ранее), мы столкнемся с физикой процессора.

Глобальная очередь становится точкой сериализации. Чтобы система масштабировалась линейно, потоки вообще не должны пересекаться в памяти при поиске работы.

Децентрализация: путь к Work Stealing

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

Когда поступает новая задача, балансировщик (или сам создающий поток) кладет ее в локальную очередь одного из воркеров. Теперь потоки не конкурируют друг с другом: каждый берет задачи только из своей очереди. Конкуренция (Contention) падает до нуля.

Но возникает новая проблема — дисбаланс нагрузки. Задачи редко бывают абсолютно одинаковыми по времени выполнения. Поток AA может получить 10 тяжелых задач и трудиться секунду, а поток BB получит 10 легких задач, выполнит их за миллисекунду и уснет. Процессор простаивает, хотя работа в системе есть.

Именно здесь на сцену выходит паттерн Work Stealing (Кража работы).

Его суть проста: если рабочий поток опустошил свою локальную очередь, он не засыпает. Вместо этого он становится «вором», заглядывает в локальные очереди других рабочих потоков и «крадет» у них задачи.

Механика Deque: элегантность без конфликтов

Чтобы Work Stealing был эффективным, нужно решить, как именно вор забирает задачу у владельца очереди, чтобы они не подрались за одни и те же ячейки памяти.

Для этого в качестве локальной очереди используется Deque (Double-Ended Queue — двусторонняя очередь). Это структура, позволяющая добавлять и извлекать элементы с обоих концов.

Паттерн Work Stealing предписывает строгие правила использования концов этой очереди:

  1. Владелец (Owner) работает со своей очередью как со стеком (LIFO — Last In, First Out). Он всегда кладет новые задачи на дно (bottom) и берет задачи для выполнения тоже со дна.
  2. Вор (Thief) всегда крадет задачи с вершины (top) чужой очереди (FIFO — First In, First Out).

Почему выбрана именно такая асимметричная схема? За этим стоят две гениальные причины, напрямую связанные с тем, что мы изучали в предыдущих главах.

1. Минимизация конкуренции (Contention)

Владелец и вор работают с разными концами очереди. Владелец меняет указатель bottom, а вор меняет указатель top. Это значит, что они модифицируют разные переменные. Если эти переменные разнесены по разным кэш-линиям (вспомните @Contended), владелец может добавлять и брать задачи вообще без синхронизации с вором. Конфликт (и необходимость дорогого CAS) возникает только в одном случае: когда в очереди осталась ровно одна задача, и владелец с вором потянулись за ней одновременно.

2. Временная локальность кэша (Temporal Locality)

Почему владелец берет самую свежую задачу (LIFO)? Потому что данные, связанные с этой задачей, скорее всего, всё ещё находятся в L1/L2 кэше ядра процессора. Если поток только что создал подзадачу, её контекст «горячий». Выполняя её немедленно, мы получаем максимальный процент попаданий в кэш (Cache Hit Ratio).

С другой стороны, вор забирает самую старую задачу с вершины. В алгоритмах «разделяй и властвуй» (Divide and Conquer) старые задачи обычно представляют собой самые крупные куски неразделенной работы. Вор забирает большой кусок, дробит его и заполняет уже свою локальную очередь, надолго обеспечивая себя работой и снижая необходимость воровать снова.

Резюме

Паттерн Work Stealing — это вершина эволюции пулов потоков для вычислительно интенсивных (CPU-bound) задач. Он объединяет децентрализацию данных (локальные очереди) для устранения узкого горлышка и динамическую балансировку (кражу) для предотвращения простоя ядер. А использование Deque с асимметричным доступом идеально ложится на архитектуру процессорных кэшей.

В следующей главе мы перейдем от теории к практике и разберем, как этот элегантный паттерн реализован внутри фреймворка ForkJoinPool, который является двигателем для Parallel Streams и CompletableFuture в Java.

Фреймворк ForkJoinPool и параллельные стримы (Parallel Streams)

Фреймворк ForkJoinPool и параллельные стримы (Parallel Streams)

Мы уже знаем, как паттерн Work Stealing решает проблему дисбаланса нагрузки: каждый поток получает собственную двустороннюю очередь (Deque) и крадет задачи у соседей, если своя очередь опустела. Но как именно большую монолитную задачу разбить на тысячи мелких, чтобы их можно было распределить по этим очередям?

В Java 7 для этого был представлен фреймворк ForkJoinPool — специализированная реализация пула потоков, построенная поверх алгоритма Work Stealing и парадигмы «разделяй и властвуй» (Divide and Conquer).

Парадигма Fork/Join

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

Алгоритм работы описывается двумя шагами:

  1. Fork (Разветвление): Если задача превышает заданный порог (threshold), она разбивается на две или более независимые подзадачи. Эти подзадачи асинхронно отправляются в пул.
  2. Join (Слияние): Текущая задача дожидается завершения своих подзадач и комбинирует их результаты.

Математически это выражается рекуррентным соотношением: время выполнения задачи T(N)T(N) сводится к параллельному выполнению подзадач плюс константное время на их слияние: T(N)=T(N/2)+O(1)T(N) = T(N/2) + O(1).

Анатомия ForkJoinTask: как работают fork() и join()

В ForkJoinPool задачи представлены абстрактным классом ForkJoinTask. У него есть две основные реализации, с которыми мы работаем:

  • RecursiveAction — для задач, не возвращающих результат (например, параллельная сортировка массива in-place).
  • RecursiveTask<V> — для задач, возвращающих значение типа V (например, вычисление суммы элементов).

Посмотрим на классический пример — поиск суммы массива.

public class SumTask extends RecursiveTask<Long> {
    private static final int THRESHOLD = 10_000;
    private final long[] array;
    private final int start;
    private final int end;

    public SumTask(long[] array, int start, int end) {
        this.array = array;
        this.start = start;
        this.end = end;
    }

    @Override
    protected Long compute() {
        if (end - start <= THRESHOLD) {
            long sum = 0;
            for (int i = start; i < end; i++) {
                sum += array[i];
            }
            return sum;
        }

        int mid = start + (end - start) / 2;
        SumTask leftTask = new SumTask(array, start, mid);
        SumTask rightTask = new SumTask(array, mid, end);

        leftTask.fork(); // Асинхронно отправляем в очередь
        long rightResult = rightTask.compute(); // Вычисляем правую часть в текущем потоке
        long leftResult = leftTask.join(); // Дожидаемся левую часть

        return leftResult + rightResult;
    }
}

Здесь кроется главная магия фреймворка, отличающая его от обычного ThreadPoolExecutor.

Когда мы вызываем leftTask.fork(), задача помещается на дно (bottom) локальной очереди текущего потока (вспоминаем LIFO-доступ владельца из предыдущей главы).

Но что происходит при вызове leftTask.join()? В классическом многопоточном программировании ожидание результата означает блокировку потока (переход в состояние WAITING на уровне ОС). Если бы ForkJoinPool блокировал поток при каждом join(), мы бы мгновенно исчерпали лимит потоков и получили Deadlock.

Вместо блокировки join() работает по принципу взаимопомощи (Helping):

  1. Поток проверяет, не лежит ли ожидаемая задача на дне его собственной очереди. Если да — он достает ее и выполняет (LIFO).
  2. Если задачи там нет (ее уже украл другой поток), текущий поток не засыпает. Он сам становится «вором» и начинает красть другие доступные задачи из очередей соседних потоков, выполняя полезную работу, пока ожидаемый результат не будет готов.

Parallel Streams: фасад для ForkJoinPool

В Java 8 появились Stream API и метод parallel(), который сделал параллельные вычисления доступными в одну строку. Под капотом любой параллельный стрим использует общий пул потоков — ForkJoinPool.commonPool().

Размер этого пула по умолчанию равен количеству логических ядер процессора минус один (одно ядро оставляется для потока, запускающего стрим — main thread).

Но чтобы ForkJoinPool мог применить парадигму Fork/Join, источник данных стрима нужно как-то разделить на части. За это отвечает интерфейс Spliterator (Split + Iterator). Метод trySplit() пытается отщепить от текущего набора данных кусок и вернуть новый Spliterator для этого куска.

Эффективность параллельного стрима критически зависит от того, насколько дешево структуре данных выполнять trySplit().

  • Массивы и ArrayList: Идеальные кандидаты. trySplit() просто вычисляет индекс середины и возвращает Spliterator для второй половины. Это операция O(1)O(1), не требующая копирования данных.
  • LinkedList: Худший кандидат. Чтобы найти середину связного списка, нужно пройти по ссылкам от начала до середины. Это операция O(N)O(N), которая убивает всю выгоду от распараллеливания.
  • HashSet / TreeSet: Разделяются сложнее, чем массивы, но лучше, чем связные списки. Разделение происходит на уровне корзин (buckets) или узлов дерева.

Антипаттерны: когда .parallel() убивает производительность

Легкость вызова .parallel() породила множество проблем в production-системах. Параллельные стримы эффективны только для CPU-bound задач на больших и легко делимых объемах данных. В остальных случаях они наносят вред.

1. Блокирующий ввод-вывод (I/O) в стриме

ForkJoinPool.commonPool() — это глобальный ресурс для всей JVM. Если внутри параллельного стрима сделать HTTP-запрос или обращение к БД, рабочий поток пула заблокируется в ожидании ответа по сети.

Поскольку размер commonPool равен количеству ядер (например, 8), всего 8 медленных HTTP-запросов полностью исчерпают пул. Все остальные параллельные стримы в приложении (даже в других модулях) выстроятся в очередь и заблокируются.

2. Зависимость от состояния (Stateful operations)

Операции вроде distinct(), sorted(), limit() или skip() требуют знания о других элементах стрима. Например, limit(10) в последовательном стриме просто остановит обработку после 10 элементов. В параллельном стриме потоки обрабатывают данные хаотично. Чтобы вернуть строго первые 10 элементов исходного массива, фреймворку приходится вводить дорогостоящую синхронизацию между потоками, что сводит на нет ускорение.

3. Накладные расходы на мелких задачах

Передача задачи в очередь, кража (Work Stealing), слияние результатов — все это требует процессорного времени. Если вычисление элементарного математического выражения над массивом из 1000 элементов занимает микросекунды, накладные расходы фреймворка превысят время самой полезной работы. Параллельные стримы начинают окупаться на объемах от десятков тысяч элементов (или если обработка одного элемента крайне ресурсоемка).

Понимание того, как ForkJoinPool дробит задачи и как потоки не блокируются благодаря Work Stealing, позволяет писать высокопроизводительный код. Однако для задач с вводом-выводом и сетевых вызовов эта архитектура не подходит.

Реактивное программирование и Flow API в Java

Реактивное программирование и Flow API в Java

В прошлой главе мы увидели, как ForkJoinPool и параллельные стримы блестяще справляются с вычислительными (CPU-bound) задачами, разделяя их на части. Но мы также столкнулись с фатальной проблемой: стоило добавить внутрь стрима HTTP-запрос или обращение к базе данных, как commonPool мгновенно исчерпывался.

Причина кроется в природе традиционного ввода-вывода (I/O). Когда поток вызывает InputStream.read(), он переходит в состояние ожидания на уровне операционной системы. Поток ничего не вычисляет, но продолжает удерживать около 1 МБ памяти под свой стек и остается в структурах данных планировщика ОС. Если у вас 10 000 одновременных сетевых соединений, вам потребуется 10 000 потоков и 10 ГБ оперативной памяти просто для того, чтобы ждать данные.

Чтобы масштабировать I/O без бесконечного создания потоков, парадигму нужно перевернуть: поток не должен ждать данных. Он должен получать управление только тогда, когда данные уже готовы к обработке.

Смена парадигмы: Pull против Push

Традиционный подход к коллекциям и потокам данных в Java основан на модели Pull (вытягивание). Потребитель запрашивает следующий элемент и блокируется, пока не получит его. Классический пример — интерфейс Iterable<T>.

Реактивное программирование меняет эту модель на Push (выталкивание). Источник данных сам уведомляет потребителя о появлении нового элемента, об ошибке или о завершении потока данных. Потребитель реагирует на эти события асинхронно.

Характеристика Pull-модель (Iterable<T>) Push-модель (Publisher<T>)
Инициатор Потребитель (вызывает next()) Производитель (вызывает onNext())
Поведение при ожидании Поток блокируется Поток освобождается (асинхронно)
Сигнал завершения hasNext() возвращает false Вызов метода onComplete()
Обработка ошибок Выброс исключения (try/catch) Вызов метода onError()

Проблема быстрого производителя и медленного потребителя

Переход на Push-модель рождает новую проблему. Что произойдет, если производитель читает данные из быстрого кэша в памяти и пушит их со скоростью 100 000 элементов в секунду, а потребитель записывает их в медленную базу данных со скоростью 1 000 элементов в секунду?

В главе 16 мы решали эту проблему с помощью паттерна Producer-Consumer и BlockingQueue. Если очередь заполнялась, мы применяли физическое Backpressure (обратное давление) — поток-производитель блокировался на мониторе или ReentrantLock.

Но в реактивном мире блокировать потоки запрещено. Нам нужен механизм асинхронного обратного давления. Потребитель должен иметь возможность сказать производителю: «Я готов принять ровно NN элементов, отправь их и остановись, пока я не попрошу еще».

Стандарт Java 9: Flow API

До Java 9 в экосистеме существовало множество реактивных библиотек (RxJava, Project Reactor, Akka Streams), каждая со своими интерфейсами. Чтобы они могли обмениваться данными друг с другом без адаптеров, в Java 9 был добавлен класс java.util.concurrent.Flow, содержащий четыре вложенных интерфейса. Это не реализация реактивного фреймворка, это строгий контракт (SPI — Service Provider Interface).

1. Publisher (Издатель)

Источник данных. Имеет всего один метод:

public interface Publisher<T> {
    void subscribe(Subscriber<? super T> subscriber);
}
2. Subscriber (Подписчик)

Приемник данных. Содержит методы для реакции на события жизненного цикла:

public interface Subscriber<T> {
    void onSubscribe(Subscription subscription);
    void onNext(T item);
    void onError(Throwable throwable);
    void onComplete();
}
3. Subscription (Подписка)

Самый важный интерфейс, связующее звено между Publisher и Subscriber. Именно он реализует асинхронное обратное давление:

public interface Subscription {
    void request(long n);
    void cancel();
}
4. Processor (Процессор)

Компонент, который одновременно является и подписчиком, и издателем. Используется для трансформации данных в пайплайне (например, map или filter).

public interface Processor<T, R> extends Subscriber<T>, Publisher<R> {
}

Жизненный цикл и математика Backpressure

Взаимодействие между этими интерфейсами подчиняется строгому протоколу, который можно описать регулярным выражением: onSubscribe onNext* (onError | onComplete)?.

  1. Потребитель вызывает publisher.subscribe(subscriber).
  2. Издатель создает объект Subscription и передает его потребителю через вызов subscriber.onSubscribe(subscription).
  3. Потребитель сохраняет подписку и запрашивает первую порцию данных: subscription.request(n).
  4. Издатель начинает вызывать subscriber.onNext(item), но не более nn раз.

Метод request(long n) — это математический контракт. У издателя есть внутренний счетчик PendingPending (изначально равен 0). При вызове request(n) происходит операция Pending=Pending+nPending = Pending + n. При каждом вызове onNext() счетчик уменьшается: Pending=Pending1Pending = Pending - 1. Издатель имеет право вызывать onNext() только при условии Pending>0Pending > 0. Если счетчик достигает нуля, издатель обязан приостановить генерацию данных, не блокируя при этом свой поток (например, отписавшись от событий нижележащего NIO-канала).

Посмотрим, как выглядит корректная реализация подписчика, обрабатывающего элементы по одному:

import java.util.concurrent.Flow.*;

public class SlowDatabaseSubscriber implements Subscriber<String> {
    private Subscription subscription;

    @Override
    public void onSubscribe(Subscription subscription) {
        this.subscription = subscription;
        // Запрашиваем только 1 элемент для старта
        this.subscription.request(1);
    }

    @Override
    public void onNext(String item) {
        // Имитация медленной асинхронной записи в БД
        saveToDatabaseAsync(item).thenRun(() -> {
            // Только после успешной обработки запрашиваем следующий элемент
            this.subscription.request(1);
        });
    }

    @Override
    public void onError(Throwable throwable) {
        System.err.println("Ошибка потока: " + throwable.getMessage());
    }

    @Override
    public void onComplete() {
        System.out.println("Все данные успешно сохранены.");
    }

    private CompletableFuture<Void> saveToDatabaseAsync(String item) {
        // ... неблокирующий I/O вызов ...
    }
}

Почему в Java только интерфейсы?

Вы можете заметить, что в стандартной библиотеке Java (до появления виртуальных потоков) практически нет готовых реализаций Publisher. Flow API задумывался как язык общения между мощными сторонними библиотеками.

На практике вы редко будете реализовывать Subscriber вручную. Вы будете использовать Project Reactor (классы Flux и Mono) или RxJava (класс Flowable). Эти библиотеки предоставляют сотни операторов (map, flatMap, zip, retry), которые под капотом виртуозно управляют вызовами request(n), объединяя потоки данных и пробрасывая обратное давление от самого медленного потребителя к самому быстрому источнику, не блокируя при этом ни одного системного потока.

Понимание контракта Flow API позволяет осознанно писать неблокирующий код и понимать, почему Flux не начнет выполнять I/O запрос до тех пор, пока на него не подпишутся и не вызовут метод request().

Поиск утечек памяти и взаимных блокировок с помощью jstack и jvisualvm

Поиск утечек памяти и взаимных блокировок с помощью jstack и jvisualvm

Сервер перестает отвечать на health-check запросы балансировщика, хотя загрузка процессора падает до нуля. Или, наоборот, потребление памяти медленно, но верно ползет вверх, пока приложение не падает с OutOfMemoryError, унося с собой необработанные транзакции. Идеально спроектированные неблокирующие алгоритмы и реактивные пайплайны бессильны, если в коде затаилась логическая ошибка удержания ресурсов. Когда система зависает в production-среде, исходный код уже не поможет — нужно уметь читать состояние работающей виртуальной машины.

Анатомия зависания: чтение Thread Dump через jstack

Когда приложение замирает, первое действие инженера — получить снимок состояния всех потоков (Thread Dump). Утилита jstack, поставляемая вместе с JDK, позволяет заглянуть внутрь процесса и увидеть, на какой именно строке кода остановился каждый поток и чего он ждет.

Вызов предельно прост: jstack <PID>, где PID — идентификатор процесса Java. Результат представляет собой текстовый файл, в котором для каждого потока выведен его стек вызовов и текущее состояние.

Ключевой навык при анализе дампа — умение мгновенно фильтровать системные потоки (сборщик мусора, JIT-компилятор) и фокусироваться на потоках приложения, обращая внимание на их статус:

  • RUNNABLE — поток выполняется на процессоре или готов к выполнению (ждет такта планировщика ОС). Если приложение зависло, а потоки RUNNABLE, ищите бесконечный цикл (Spin-wait) или блокирующий I/O вызов (например, чтение из сокета без таймаута).
  • BLOCKED — поток ожидает захвата монитора (synchronized), который в данный момент удерживается другим потоком.
  • WAITING / TIMED_WAITING — поток добровольно приостановил работу, ожидая сигнала. Это происходит при вызове Object.wait(), Thread.join(), LockSupport.park() (который используется под капотом ReentrantLock и CompletableFuture).

Если в системе произошел классический Deadlock на мониторах, jstack автоматически обнаружит цикл графа ожиданий и выведет в конце отчета секцию Found one Java-level deadlock, явно указав виновников. Однако современные системы редко используют synchronized. Если взаимная блокировка произошла на объектах ReentrantLock или из-за циклического ожидания результатов CompletableFuture, автоматического детектора может не хватить — придется распутывать клубок состояний WAITING вручную, сопоставляя идентификаторы потоков и адреса объектов в памяти.

Визуализация потоков в JVisualVM

Анализ текстового дампа jstack незаменим для точечной диагностики, но он дает лишь статический срез в одну миллисекунду. Чтобы увидеть картину в динамике — кто сколько времени работает, а кто простаивает — используется JVisualVM.

Вкладка Threads в JVisualVM рисует временную шкалу (таймлайн) для каждого потока. Цветовая кодировка позволяет с одного взгляда оценить здоровье пула потоков:

  • Зеленый (Running) — полезная работа.
  • Желтый (Wait) — ожидание задач (нормально для пула) или ожидание сигнала (потенциальная проблема, если длится долго).
  • Красный (Monitor) — потоки выстроились в очередь за synchronized ресурсом (узкое горлышко).

Если вы видите, что пул потоков обработки HTTP-запросов (например, http-nio-exec-*) постоянно перекрашивается в желтый цвет, это означает, что потоки простаивают, ожидая ответов от базы данных или стороннего API. Это классический сигнал к тому, что синхронная модель исчерпала себя и пора переходить к асинхронному I/O или Flow API.

Специфика утечек памяти в многопоточной среде

В языках с ручным управлением памятью утечка — это забытый вызов free(). В Java, благодаря сборщику мусора (GC), классических утечек не бывает. Здесь утечка памяти — это непреднамеренное удержание ссылок (Unintentional Object Retention). Объект больше не нужен логике приложения, но ссылка на него все еще существует в «живой» части графа объектов, поэтому GC не имеет права его удалить.

Многопоточность создает два специфических и очень коварных паттерна утечек памяти.

1. Бесконечный рост очередей в Thread Pool

Паттерн Producer-Consumer часто реализуется через ExecutorService с очередью задач. Если производители (Producers) генерируют задачи быстрее, чем рабочие потоки (Consumers) успевают их обрабатывать, очередь начинает расти. Поскольку задачи часто содержат тяжеловесные контексты (например, загруженные из БД сущности), память быстро исчерпывается. В логах это выглядит парадоксально: приложение активно работает, CPU загружен, но в итоге падает с OutOfMemoryError: Java heap space. Проблема решается ограничением размера очереди (Bounded Queue) и настройкой политики отказа (например, CallerRunsPolicy).

2. Отравление ThreadLocal в пулах потоков

Переменные ThreadLocal привязывают данные к конкретному потоку. Значения хранятся в специальной структуре ThreadLocalMap, которая физически является полем объекта Thread.

Пока поток жив, жива и его карта. В простых приложениях поток завершается, и GC очищает его память. Но в серверных приложениях потоки переиспользуются пулом (Thread Pool). Поток живет неделями. Если в процессе обработки запроса вы положили данные в ThreadLocal и не вызвали remove() в блоке finally, эти данные останутся висеть в памяти навсегда. При следующем запросе этот же поток может перезаписать значение, но если ключ больше не используется, старый объект останется в памяти, накапливаясь с каждой итерацией.

Поиск корня зла: анализ Heap Dump

Когда JVisualVM показывает ступенчатый рост потребления кучи (Heap), который не сбрасывается после вызова GC, необходимо сделать снимок памяти (Heap Dump). Это можно сделать прямо из интерфейса JVisualVM кнопкой «Heap Dump» или утилитой командной строки jmap -dump:live,format=b,file=heap.bin <PID>.

Анализ дампа (через тот же JVisualVM или Eclipse MAT) сводится к поиску объектов, потребляющих больше всего памяти. При этом критически важно различать два понятия:

  • Shallow Size — размер самого объекта (заголовки, примитивы, ссылки). Обычно это десятки байт.
  • Retained Size — объем памяти, который будет освобожден сборщиком мусора, если удалить этот конкретный объект. Это сумма Shallow Size самого объекта и всех объектов, на которые ссылается только он.

Если вы нашли объект с аномально большим Retained Size (например, массив Object[] внутри LinkedBlockingQueue), следующий шаг — найти GC Root. GC Root — это точка входа в граф объектов (локальная переменная выполняющегося метода, статический класс или активный поток).

Инструмент анализа позволяет построить путь от утекшего объекта к GC Root (Path to GC Roots). Именно этот путь показывает, кто держит объект. Если путь ведет к объекту java.lang.Thread, а от него к ThreadLocalMap, вы точно знаете: проблема в забытом вызове remove(). Если путь ведет к ThreadPoolExecutor, проблема в переполнении очереди задач.

Умение сопоставлять статические снимки памяти (Heap Dump) с динамикой работы потоков (Thread Dump) позволяет находить архитектурные просчеты, которые не выявляются модульными тестами, и обеспечивает стабильность системы под реальной нагрузкой.

Профилирование многопоточного кода с JMH (Java Microbenchmark Harness)

Профилирование многопоточного кода с JMH (Java Microbenchmark Harness)

В прошлой главе мы научились находить утечки памяти и взаимные блокировки с помощью дампов. Но когда код работает корректно, возникает следующий вопрос: насколько он быстр?

Типичный подход начинающего разработчика — замерить время до выполнения кода через System.nanoTime(), запустить код в цикле и замерить время после. В мире Java этот подход гарантированно выдаст ложные результаты. JIT-компилятор (о котором мы говорили в главе про Escape Analysis) настолько умен, что может просто удалить ваш код, если поймет, что его результаты нигде не используются. Операционная система может поставить поток на паузу, а сборщик мусора — заморозить приложение (Stop-The-World) прямо посреди замера.

Чтобы измерять производительность с хирургической точностью, инженеры OpenJDK создали JMH (Java Microbenchmark Harness) — стандартный фреймворк для написания микробенчмарков на уровне JVM.

Анатомия обмана: почему наивные тесты лгут

Представьте, что вы хотите замерить скорость инкремента AtomicInteger. Вы пишете цикл на миллион итераций. Что сделает JIT-компилятор?

  1. Warmup (Прогрев). Первые тысячи вызовов код интерпретируется. Он медленный. Наивный тест включит это время в общий результат, хотя в реальном production-приложении этот код давно был бы скомпилирован в машинные инструкции.
  2. Dead Code Elimination (DCE). Если вы просто вызываете atomic.incrementAndGet() и никуда не сохраняете результат, JIT-компилятор (на фазе агрессивных оптимизаций) поймет, что операция не имеет наблюдаемых побочных эффектов вне цикла. Он может схлопнуть весь цикл в одну математическую операцию или вообще удалить его.
  3. Loop Unrolling (Развертывание циклов). Накладные расходы на проверку условия цикла (i<1000000i < 1000000) могут превысить время самой полезной работы внутри него.

JMH решает эти проблемы системно. Он самостоятельно организует фазы прогрева (Warmup) и измерения (Measurement), запускает код в изолированных форках JVM (чтобы профили старых тестов не влияли на новые) и предоставляет инструменты для обмана JIT-компилятора.

Разделение состояния: @State в многопоточной среде

В многопоточном бенчмарке критически важно отделить данные, к которым потоки обращаются конкурентно (разделяемое состояние), от данных, которые принадлежат только одному потоку (локальное состояние). В JMH за это отвечает аннотация @State.

Существует две главные области видимости (Scope):

  • @State(Scope.Benchmark) — один экземпляр объекта создается на весь запуск бенчмарка. Все рабочие потоки будут обращаться к этому единственному объекту. Это идеальное место для размещения нашей тестируемой структуры данных (например, ConcurrentHashMap или LongAdder).
  • @State(Scope.Thread) — каждый рабочий поток получит свою собственную, независимую копию объекта. Это нужно для генераторов случайных чисел, локальных буферов или счетчиков, чтобы потоки не мешали друг другу при подготовке тестовых данных.

Рассмотрим пример подготовки бенчмарка для сравнения AtomicLong и LongAdder (вспоминаем главу 20):

@State(Scope.Benchmark)
public class CounterState {
    public AtomicLong atomic = new AtomicLong();
    public LongAdder adder = new LongAdder();
}

@State(Scope.Thread)
public class ThreadLocalState {
    // У каждого потока будет свой генератор,
    // чтобы избежать contention на самом генераторе
    public ThreadLocalRandom random = ThreadLocalRandom.current();
}

Защита от оптимизаций: Blackhole

Чтобы JIT-компилятор не удалил наш код через Dead Code Elimination, мы должны убедить его, что результат вычислений критически важен. В JMH есть два способа сделать это.

Первый — возвращать результат прямо из метода с аннотацией @Benchmark:

@Benchmark
public long measureAtomic(CounterState state) {
    return state.atomic.incrementAndGet();
}

JMH неявно передаст это возвращаемое значение в специальный поглотитель, и JIT будет обязан честно выполнить инструкцию lock xadd (CAS).

Второй способ — использовать объект Blackhole, который JMH передает в параметры метода. Это необходимо, когда нужно «поглотить» несколько значений или когда метод возвращает void:

@Benchmark
public void measureMultiple(CounterState state, Blackhole bh) {
    bh.consume(state.atomic.incrementAndGet());
    bh.consume(state.atomic.get());
}

Blackhole спроектирован так, чтобы JIT не мог предсказать его поведение. Передача переменной в черную дыру гарантирует, что код, вычисляющий эту переменную, будет выполнен физически.

Асимметричные нагрузки: @Group и @GroupThreads

Многие многопоточные структуры данных ведут себя по-разному в зависимости от профиля нагрузки. Например, ReentrantReadWriteLock или StampedLock показывают чудеса производительности, когда 99% потоков только читают данные, и катастрофически деградируют, если половина потоков начинает писать.

Просто запустить 10 потоков, каждый из которых делает случайный выбор (читать или писать), — плохая идея. Это создаст непредсказуемый профиль нагрузки, зависящий от случайности в конкретную миллисекунду.

JMH позволяет жестко задать роли потоков с помощью аннотаций @Group и @GroupThreads.

@State(Scope.Benchmark)
public class CacheState {
    private final Map<String, String> cache = new ConcurrentHashMap<>();

    @Setup
    public void setup() {
        cache.put("key", "value");
    }

    public String read() { return cache.get("key"); }
    public void write(String val) { cache.put("key", val); }
}

@Benchmark
@Group("cacheTest")
@GroupThreads(3) // 3 потока будут непрерывно читать
public String reader(CacheState state) {
    return state.read();
}

@Benchmark
@Group("cacheTest")
@GroupThreads(1) // 1 поток будет непрерывно писать
public void writer(CacheState state, ThreadLocalState ts) {
    state.write(String.valueOf(ts.random.nextInt()));
}

В этом сценарии JMH создаст группу cacheTest из 4 потоков. Три из них будут бесконечно вызывать метод reader, а один — метод writer. JMH синхронизирует их запуск, соберет статистику (пропускную способность) для каждого метода отдельно, а затем выдаст агрегированный результат для всей группы.

Метрики: что именно мы измеряем?

По умолчанию JMH измеряет Throughput (пропускную способность) — количество операций в единицу времени. Формула проста: Throughput=OperationsTimeThroughput = \frac{Operations}{Time}. Результат выдается в операциях в секунду (ops/s). Это главная метрика для оценки масштабируемости lock-free алгоритмов.

Однако для систем реального времени или пользовательских интерфейсов важнее задержка. В JMH можно изменить режим через аннотацию @BenchmarkMode:

  • Mode.AverageTime — среднее время выполнения одной операции (например, наносекунд на вызов).
  • Mode.SampleTime — собирает распределение времени выполнения (перцентили: 90%, 99%, 99.9%). Это критически важно для выявления редких, но долгих пауз (например, когда поток попадает в длинный цикл spin-wait или ждет сборщика мусора).

Профилирование многопоточного кода — это наука о нюансах. JMH дает вам микроскоп, но интерпретация увиденного остается за инженером. Увидели высокую пропускную способность, но плохой 99-й перцентиль? Возможно, ваш алгоритм страдает от голодания (Starvation) отдельных потоков. Увидели падение throughput при увеличении числа потоков? Ищите ложное разделение (False Sharing) или избыточный contention на CAS-операциях.

В следующей, заключительной главе мы объединим все наши знания — от атомиков и JMM до профилирования и пулов потоков — чтобы спроектировать архитектуру реального отказоустойчивого сервиса обработки платежей.

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

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

В Черную пятницу шлюз приема платежей получает 10 000 запросов в секунду. Внешний банковский API, который обычно отвечает за 50 миллисекунд, внезапно начинает «тормозить» и отвечать за 2 секунды. Если ваш сервис написан по принципу «один поток на один запрос» с синхронным ожиданием ответа, пул потоков сервера (например, Tomcat) исчерпается ровно за 0.2 секунды. Сервис перестанет принимать новые соединения, включая те, что не требуют обращения к банку. Система упадет.

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

Изоляция сбоев: паттерн Bulkhead

Фундаментальная проблема синхронной обработки — каскадные сбои. Если все задачи выполняются в едином пуле потоков (или ForkJoinPool.commonPool()), медленный внешний ресурс захватит все доступные потоки. Остальные компоненты системы окажутся заблокированы в ожидании свободных воркеров (Thread Starvation).

Решение пришло из кораблестроения — паттерн Bulkhead (переборки). Корпус судна делится на изолированные отсеки: если один пробит, вода не затапливает остальные. В многопоточном программировании «отсеки» — это физически разделенные пулы потоков (Thread Pools) для разных типов задач.

В платежном сервисе мы выделяем минимум три пула:

  1. Acceptor Pool (I/O потоки сервера) — только принимает HTTP-запросы, выполняет быструю валидацию и немедленно передает задачу дальше.
  2. Database Pool — потоки, настроенные под размер пула соединений с БД (например, HikariCP). Их количество жестко ограничено, чтобы не создавать конкуренцию на уровне базы данных.
  3. Bank API Pool — потоки для сетевых вызовов к внешнему шлюзу.

Если банковский API замедляется, очередь задач Bank API Pool заполняется, и он начинает отклонять новые задачи (срабатывает RejectedExecutionHandler). При этом Acceptor Pool остается свободным и может мгновенно возвращать клиентам статус HTTP 503 Service Unavailable или ставить платежи в персистентную очередь, не обрушая весь сервис.

Идемпотентность и дедупликация в памяти

При сетевых сбоях клиент (или мобильное приложение) неизбежно выполнит повторную попытку (Retry). Сервис может получить два HTTP-запроса на списание средств по одному и тому же заказу одновременно. Мы обязаны гарантировать идемпотентность — сколько бы раз ни пришел запрос с одним и тем же ключом (Idempotency Key), деньги должны списаться ровно один раз.

Классическая проверка в базе данных через SELECT перед INSERT уязвима для состояния гонки (Data Race). Использование уникальных констрейнтов в БД решает проблему, но создает высокую нагрузку на диск при массированных повторах.

Эффективный подход — слой дедупликации в оперативной памяти на базе ConcurrentHashMap.

private final ConcurrentHashMap<String, Payment> activePayments = new ConcurrentHashMap<>();

public Payment getOrCreatePayment(String idempotencyKey, BigDecimal amount) {
    return activePayments.computeIfAbsent(idempotencyKey, key -> {
        // Эта лямбда выполнится строго один раз для уникального ключа
        return new Payment(key, amount);
    });
}

Метод computeIfAbsent гарантирует атомарность: если два потока одновременно принесут одинаковый idempotencyKey, объект Payment будет создан только одним потоком. Второй поток заблокируется на уровне сегмента (или корзины) хэш-таблицы на микросекунды и получит уже созданный экземпляр.

Конечный автомат на базе CAS

Получив единственный экземпляр Payment, мы должны провести его через жизненный цикл: NEW \rightarrow PROCESSING \rightarrow SUCCESS (или FAILED).

Поскольку разные этапы обработки могут выполняться разными потоками (из-за паттерна Bulkhead), состояние платежа должно обновляться потокобезопасно. Использование блокировок (synchronized) вокруг объекта платежа снизит пропускную способность. Вместо этого мы реализуем неблокирующий конечный автомат (State Machine) с помощью AtomicReference и операции Compare-And-Swap (CAS).

public class Payment {
    private final String id;
    private final AtomicReference<State> state = new AtomicReference<>(State.NEW);

    public boolean tryTransitionToProcessing() {
        return state.compareAndSet(State.NEW, State.PROCESSING);
    }

    public void complete(State finalState) {
        // Платеж может перейти в финал только из состояния PROCESSING
        if (!state.compareAndSet(State.PROCESSING, finalState)) {
            throw new IllegalStateException("Неверный переход состояния");
        }
    }
}

Если два потока-дубликата прошли слой дедупликации (например, первый запрос уже был в обработке, а второй только пришел), вызов tryTransitionToProcessing() разрешит конфликт. Только один поток получит true и отправит запрос в банк. Второй получит false и сможет просто подписаться на результат работы первого.

Асинхронная оркестрация через CompletableFuture

Теперь необходимо связать пулы потоков и конечный автомат воедино, не блокируя потоки на ожидании. Для этого используется CompletableFuture. Каждый этап конвейера возвращает фьючерс, а переходы между пулами осуществляются через методы с суффиксом *Async.

public CompletableFuture<PaymentResult> process(PaymentRequest req) {
    Payment payment = getOrCreatePayment(req.getIdempotencyKey(), req.getAmount());

    if (!payment.tryTransitionToProcessing()) {
        // Платеж уже обрабатывается другим потоком, возвращаем существующий пайплайн
        return payment.getResultFuture();
    }

    CompletableFuture<PaymentResult> pipeline = CompletableFuture
        // 1. Сохраняем в БД (используем dbPool)
        .supplyAsync(() -> repository.save(payment), dbPool)

        // 2. Отправляем в банк (переключаемся на bankApiPool)
        .thenComposeAsync(saved -> bankClient.charge(saved), bankApiPool)

        // 3. Обрабатываем результат (в пуле по умолчанию или вызывающем)
        .handle((bankResponse, exception) -> {
            if (exception != null) {
                payment.complete(State.FAILED);
                return PaymentResult.fail(exception);
            }
            payment.complete(State.SUCCESS);
            activePayments.remove(req.getIdempotencyKey()); // Очищаем кэш
            return PaymentResult.success(bankResponse);
        });

    payment.setResultFuture(pipeline);
    return pipeline;
}

В этой архитектуре ни один поток не простаивает в ожидании I/O. Поток из dbPool сохраняет данные и немедленно берет следующую задачу. Сетевой вызов делегируется в bankApiPool. Когда банк отвечает, коллбэк handle финализирует состояние через CAS.

Неблокирующий Exponential Backoff

Что делать, если банковский API вернул ошибку HTTP 429 Too Many Requests? Нам нужно повторить запрос через некоторое время.

Наивный подход — использовать Thread.sleep(1000) внутри bankApiPool — катастрофичен. Это мгновенно исчерпает пул. Поток будет физически заблокирован, не выполняя полезной работы.

Правильный паттерн — асинхронное отложенное выполнение. Мы используем ScheduledExecutorService только для отсчета времени, а саму работу возвращаем обратно в пул банковских запросов.

private CompletableFuture<BankResponse> retryAsync(Payment payment, int attempt) {
    if (attempt > MAX_RETRIES) return CompletableFuture.failedFuture(new TimeoutException());

    CompletableFuture<BankResponse> future = new CompletableFuture<>();
    long delayMs = (long) Math.pow(2, attempt) * 100; // 100ms, 200ms, 400ms...

    // Планируем выполнение в будущем БЕЗ блокировки текущего потока
    scheduler.schedule(() -> {
        bankClient.charge(payment)
            .thenAccept(future::complete)
            .exceptionally(ex -> {
                // Рекурсивный асинхронный вызов
                retryAsync(payment, attempt + 1).thenAccept(future::complete);
                return null;
            });
    }, delayMs, TimeUnit.MILLISECONDS);

    return future;
}

Таймер scheduler использует всего один поток, который управляет очередью тайм-аутов (внутри это структура DelayQueue). Он не выполняет сетевые запросы, а лишь в нужный момент времени кладет задачу обратно в bankApiPool. Это позволяет сервису «держать в уме» сотни тысяч отложенных платежей, потребляя минимум ресурсов процессора и памяти.

Итог

Отказоустойчивая многопоточная система не рождается из хаотичного разбрасывания synchronized или создания потоков new Thread(). Это строгий инженерный расчет:

  1. Разделение ресурсов (Bulkhead) предотвращает каскадные сбои.
  2. Неблокирующие структуры (ConcurrentHashMap) защищают от дубликатов на входе.
  3. Атомарные переменные (AtomicReference + CAS) обеспечивают консистентность состояния без блокировок.
  4. Реактивные пайплайны (CompletableFuture) утилизируют CPU на 100%, не позволяя потокам простаивать в ожидании I/O.

Именно такой подход позволяет современным Java-приложениям обрабатывать десятки тысяч транзакций в секунду на стандартном оборудовании.