Цели главы. Собрать из уже изученных деталей самый влиятельный архитектурный примитив последнего десятилетия — распределённый журнал событий. В Kafka читатель курса узнает старых знакомых: партиции — это шардирование (глава 8), реплики партиции — лидерская репликация с выбором acks (глава 5), контроллер — Raft (глава 6), а семантики доставки — теоремы главы 1 в производственной упаковке. Разберём модель, гарантии и их мелкий шрифт, честно препарируем «exactly-once» и посмотрим на журнал как на архитектурный стиль. Лабораторная — потеря сообщений и дубликаты своими руками.
Асинхронная связь через брокер (мотивация — в статье о надёжности: развязка отказов, сглаживание пиков) существует в двух жанрах:
Словарь минимален. Топик — именованный поток событий. Топик разрезан на партиции — и это в точности шардирование главы 8: партиция выбирается хешем ключа сообщения, ключ проектируется по тем же правилам (события одного заказа — с ключом заказа — попадают в одну партицию). Порядок гарантирован только внутри партиции — первое и главное, что надо помнить при проектировании топиков; между партициями порядка нет, и число партиций задаёт потолок параллелизма потребления.
Каждая партиция реплицирована: лидер принимает записи, ведомые копируют — лидерская репликация главы 5 под именем ISR (in-sync replicas — множество не отставших реплик). Метаданные кластера и выбор лидеров партиций ведёт контроллер — с современных версий это встроенный Raft (режим KRaft, сменивший ZooKeeper): глава 6 в титрах. Группы потребителей: партиции топика распределяются между членами группы (каждая партиция — ровно одному члену); прогресс группы — зафиксированные оффсеты, хранимые в самой Kafka; изменение состава группы вызывает ребаланс — перераздачу партиций, о цене которого — в 10.4.
Рис. 10.1. Партиция является одновременно единицей порядка, репликации и назначения участнику группы.
Производитель выбирает уровень подтверждения — и это дословно выбор sync/async из главы 5 со всеми его последствиями:
Разумная исходная политика для критичных данных: acks=all, min.insync.replicas ≥ 2, replication.factor ≥ 3 и unclean.leader.election.enable=false. Но конкретные числа следуют из требуемого числа переносимых отказов и размещения реплик: сама по себе эта комбинация не защищает от потери всех копий или общей зоны отказа.
Применим главу 1. Производитель: не получив подтверждения, он повторяет отправку — брокер может получить дубликат (at-least-once). Лечение встроено: идемпотентный производитель (идентификатор производителя + порядковые номера; брокер отбрасывает повторы) снимает дубли, вызванные повторными отправками в рамках поддерживаемой сессии. В актуальных клиентах Kafka идемпотентность обычно включена по умолчанию, если настройки ей не противоречат; это не дедупликация произвольных повторных бизнес-событий.
Потребитель: вечная дилемма момента фиксации оффсета. Зафиксировал до обработки — падение теряет сообщение (at-most-once); после обработки — падение между обработкой и фиксацией даёт повторную обработку (at-least-once). Ребаланс группы устраивает то же самое массово: партиция переезжает к соседу, начинающему с последнего зафиксированного оффсета — всё после него обрабатывается повторно. Вывод неотменим: потребитель обязан быть идемпотентным — дедупликация по ключу события (глава 9, outbox: идентификатор события — ключ идемпотентности) либо атомарная фиксация «результат + оффсет» в одном хранилище (оффсет — строка в той же БД-транзакции, что и результат обработки).
А что же знаменитый exactly-once Kafka? Читаем честно: транзакции Kafka атомарно связывают «прочитанные оффсеты + записанные сообщения» внутри конвейера Kafka→Kafka (потребители читают только зафиксированные транзакции) — это exactly-once processing для потоковых конвейеров (Kafka Streams), и в своих рамках он работает. Он не отменяет двух генералов: доставка во внешний мир (запись в стороннюю БД, вызов стороннего API, письмо пользователю) остаётся at-least-once + идемпотентность — и любой лозунг «ровно один раз конец-в-конец» должен предъявить, где именно спрятана дедупликация.
Перечитываемость превращает журнал из транспорта в источник истины:
Трезвая оговорка в духе курса: «журнал всего» — не бесплатная архитектура (эксплуатация кластера, дисциплина схем, глубина удержания — деньги), и для системы из трёх сервисов outbox плюс обычная очередь может быть честнее полновесного event sourcing. Распределяйтесь — и журналируйте — от необходимости.
Ответы и указания. 1: ключ — идентификатор заказа; идентификаторы отдельных событий могут отправить события заказа в разные партиции. Без ключа современный стандартный partitioner обычно держится за одну партицию на время пакета, затем меняет её: это не обязательный round-robin и порядка конкретного заказа всё равно не гарантирует. 2: десять записей выше подтверждённой репликами границы могут исчезнуть при выборе нового лидера — аналог асинхронного хвоста PostgreSQL. При acks=all лидер ждёт все текущие ISR, а min.insync.replicas=2 запрещает успешную запись, если ISR меньше двух; с запрещённым unclean leader election подтверждённый хвост переносит отказ одного брокера. 3: фиксация оффсета до эффекта создаёт потерю при падении между ними, после эффекта — повтор. В локальной транзакции потребитель блокирует свою строку оффсета, проверяет следующий номер, записывает эффект и новый оффсет вместе; после рестарта он начинает с этой строки, а Kafka-offset в такой схеме служит не источником истины. 4: session timeout исключает участника, и новые владельцы повторят сообщения после последних зафиксированных оффсетов. Возобновившийся старый участник ограждён поколением группы от фиксации оффсета, но до следующего обращения к координатору способен повторить внешний эффект, если приложение не идемпотентно. Это детектор отказов на тайм-ауте, а не доказательство смерти. 5: транзакции Kafka связывают оффсеты и выходные записи внутри Kafka; вебхук находится за этой границей. Потерянный ответ вебхука снова создаёт неопределённость двух генералов, поэтому получателю всё ещё нужен ключ идемпотентности. 6: очередь непосредственно переотдаёт неподтверждённое задание другому воркеру и удаляет завершённое; журнал удобнее для независимых повторных обработок и месячного replay, но требует управления оффсетами и идемпотентностью. Для одноразового распределения работ классическая очередь обычно проще.
Стенд: три брокера Kafka в режиме KRaft — docker-compose.yml и команды запуска. Инструменты kafka-topics.sh, kafka-console-producer.sh и kafka-console-consumer.sh входят в образ; для управляемого эксперимента с моментом отказа удобен также программный клиент.
Задание 1. Порядок и партиции. Создайте топик orders с 3 партициями (replication-factor 3). Отправьте по 10 нумерованных событий для трёх «заказов» с ключом заказа; прочитайте топик и убедитесь: внутри заказа порядок строгий, глобального порядка нет. Повторите без ключей и запишите фактическое распределение. Короткая серия может целиком попасть в одну sticky-партицию — это нормальный результат, а не доказательство глобального порядка. Для контролируемого сравнения повторите опыт с RoundRobinPartitioner либо явно задавайте партиции программным клиентом: события одного заказа окажутся в разных независимых последовательностях, и порядка между ними Kafka не гарантирует. Требуется показать потерю гарантии, а наблюдаемая инверсия в коротком чтении необязательна. В отчёт: распределение сообщений по партициям во всех вариантах и объяснение границы гарантии.
Задание 2. Риск потери подтверждённого при acks=1. Производитель-скрипт шлёт длинную нумерованную серию с acks=1 и enable.idempotence=false, журналируя завершение каждого send; сразу после заранее выбранного подтверждения «убейте» фактического лидера исследуемой партиции, найденного через kafka-topics.sh --describe. После переизбрания сравните журнал подтверждений с содержимым топика. Реплики могут успеть скопировать весь подтверждённый хвост, поэтому нулевая потеря в отдельном запуске допустима: повторите опыт несколько раз, а при наличии средства сетевой инъекции отдельно задержите трафик репликации, не нарушая работу контроллерного кворума. Затем повторите с acks=all и min.insync.replicas=2; при сохранении хотя бы одной подтверждавшей ISR отказ одного лидера не теряет подтверждённые записи. В отчёт: номера потерянных сообщений или честный нулевой результат, точный момент отказа и условия опыта; таблица «acks → потери → задержка» (задержку измерьте на 10 000 сообщений).
Задание 3. Дубликаты при ребалансе. Группа из двух потребителей пишет обработанные номера в файлы; убейте одного потребителя между обработкой и фиксацией оффсета (для наглядности — ручная фиксация с задержкой). Пересчитайте объединение файлов: найдите повторно обработанные сообщения. В отчёт: сколько дублей и от чего зависит их число (интервал фиксации).
Задание 4. Идемпотентный потребитель. Добавьте потребителю дедупликацию: таблица обработанных идентификаторов (SQLite достаточно), проверка перед побочным эффектом, атомарная запись «результат + идентификатор». Повторите сценарий задания 3 — дублей в результате нет. В отчёт: где именно (в какой строке) at-least-once превратился в «эффект ровно один раз».
Задание 5*. Меньшинство ISR. С min.insync.replicas=2 остановите два узла из трёх. Запись должна перестать успешно завершаться, но зафиксируйте фактический результат: это может быть NOT_ENOUGH_REPLICAS, отсутствие лидера или тайм-аут. В данном компактном стенде каждый узел совмещает роли broker и controller, поэтому вместе с большинством ISR теряется и кворум KRaft — обещать один конкретный код ошибки нельзя. Что происходит с чтением? Сформулируйте поведение в терминах CAP и сравните с заданием 3 лабораторной 2 (etcd): в чём выбор Kafka совпадает и в чём даёт больше ручек настройки.