Kafka Lab

Тренажёр тестирования событийных систем на примере Apache Kafka: топики, партиции, порядок и гарантии доставки, consumer groups, дубликаты и идемпотентность, DLQ, lag и совместимость схем. Сквозной пример лабы — цепочка событий оформления заказа: OrderCreated → PaymentCaptured → OrderConfirmed → NotificationSent. Все данные учебные, к реальным брокерам/проду доступа не требуется.

Теория перед практикой · 10 разделов

Kafka — это распределённый брокер сообщений (журнал событий). Producer публикует сообщения в topic; consumer их читает. Архитектура событийная: сервисы не дёргают друг друга напрямую, а обмениваются событиями через топики, что развязывает их и позволяет масштабировать. В нашем сценарии сервис заказов публикует OrderCreated, платёжный сервис реагирует и публикует PaymentCaptured, и так по цепочке до NotificationSent.

Топик делится на partitions (партиции) — это единицы параллелизма и хранения. Каждое сообщение в партиции получает порядковый номер — offset. Consumer хранит свой offset (до какого места дочитал), и при перезапуске продолжает с него. Несколько consumer объединяются в consumer group: брокер распределяет партиции между участниками так, что каждую партицию читает ровно один член группы — так масштабируют чтение без дублирования внутри группы.

Порядок гарантируется ТОЛЬКО внутри одной партиции, глобального порядка по топику нет. Если связанные события (например, все события одного orderId) должны идти строго по порядку, их направляют в одну партицию через ключ партиционирования (partition key = orderId). Сообщения с одним ключом всегда попадают в одну партицию и читаются по порядку; без ключа Kafka раскидывает их по партициям, и порядок между ними не гарантирован — частый источник дефектов «PaymentCaptured обработался раньше OrderCreated».

Гарантии доставки (концептуально). At-most-once: сообщение может потеряться, но не задвоится (commit offset до обработки). At-least-once: сообщение точно дойдёт, но при ретраях может прийти повторно — возможны дубликаты (commit после обработки); это наиболее частый режим. Exactly-once: ровно один эффект, достигается транзакциями/идемпотентным продьюсером и стоит дороже. QA должен знать, какой режим заявлен, и тестировать его реальное поведение.

Ретраи → дубликаты → идемпотентность. При сетевом сбое producer/брокер повторяет отправку, и consumer может получить одно и то же событие дважды. Защита — идемпотентная обработка: consumer хранит обработанные eventId и при повторе не выполняет эффект второй раз. Ключевой тест нашего сценария: повторный PaymentCaptured с тем же eventId/paymentId НЕ должен списать деньги дважды и не создаст второй платёж — потребитель обязан распознать дубль по идентификатору и пропустить его.

DLQ (Dead Letter Queue) — отдельный топик для «ядовитых» сообщений, которые consumer не смог обработать после N ретраев (битый payload, несовместимая схема, бизнес-ошибка). Вместо бесконечного зацикливания на одном сообщении (которое блокирует партицию) его отправляют в DLQ для разбора. QA проверяет: попадает ли необрабатываемое событие в DLQ, сохраняется ли причина/исходный payload, не теряются ли при этом остальные сообщения, и есть ли процесс переобработки.

Consumer lag — это разница между последним записанным offset в партиции и offset, до которого дочитал consumer. Растущий lag означает, что потребитель не успевает за продьюсером: события обрабатываются с задержкой (пользователь оформил заказ, а уведомление придёт через 10 минут). Lag — ключевая метрика здоровья event-driven системы; QA проверяет lag под нагрузкой, при ребалансировке группы и после добавления/падения consumer.

Совместимость схем (schema compatibility). Структура события описана схемой (например, в Schema Registry). BACKWARD-совместимость: новые consumer читают старые сообщения (можно безопасно добавлять опциональные поля, удалять — осторожно). FORWARD: старые consumer читают новые сообщения. Несовместимое изменение (переименование/удаление обязательного поля, смена типа) ломает потребителей и часто гонит сообщения в DLQ. QA тестирует продьюсер новой версии против старого consumer и наоборот.

Observability событийных систем. Поскольку обработка асинхронна, нельзя «нажать кнопку и сразу увидеть результат». QA смотрит: дошло ли событие до топика, обработал ли его consumer (по логам с eventId), не вырос ли lag, не ушло ли что-то в DLQ, и появился ли итоговый эффект (заказ подтверждён, уведомление отправлено). Доказательство в event-driven мире — это связка eventId/correlationId через все сервисы цепочки.

Учебный payload события PaymentCaptured (JSON): {"eventId":"evt_5b21","eventType":"PaymentCaptured","paymentId":"pay_3001","orderId":"ord_9001","amount":1990,"currency":"RUB","occurredAt":"2026-06-08T10:15:42Z"}. eventId — уникальный идентификатор события (по нему ловят дубли), eventType — тип, paymentId/orderId — бизнес-ключи, occurredAt — время события. На этот payload опираются задачи про порядок, дубликаты, ключ партиционирования и DLQ.

Задача 1 / 11

Где гарантируется порядок

Порядок сообщений в Kafka — частый источник недопонимания.

Задание. Где Kafka гарантирует порядок сообщений?