← назад к разделу

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

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

4 события одного заказа order-12345 идут в один топик producer send("order-events", key=null) send("order-events", key=orderId) топик order-events · 3 партиции P0 P1 P2 consumer group consumer-1 consumer-2 payment-service shippeddeliveredпустоcreatedpaid пустоcreatedpaidshippeddeliveredпусто consumer-1consumer-2 consumer-1 порядок обработки потребители ещё не читали →→→shippedcreateddeliveredpaid →→→createdpaidshippeddelivered ключа нет — продьюсер набрал пачку в P2, потом переключился на P0 два потребителя читают параллельно: shipped раньше, чем paid ключ = orderId: hash("order-12345") % 3 = 1 → все четыре в P1 одна партиция — один потребитель: порядок как при отправке

Ключ решает, в какую партицию попадёт событие. Без ключа четыре события одного заказа разъехались по двум партициям, их читают два разных потребителя — и shipped обрабатывается раньше paid. С ключом orderId все четыре лежат в P1, а партицию в группе читает ровно один потребитель — порядок сохраняется.

Обязательно

Проблема: сервисы связаны напрямую

Представьте: пользователь оформляет заказ. OrderService должен уведомить PaymentService, InventoryService и NotificationService. Самый простой вариант — HTTP-вызовы напрямую.

Но у этого подхода есть проблема: если PaymentService недоступен, OrderService тоже не может завершить работу. Один падающий сервис тянет за собой другие. Это называют синхронной связностью — сервисы зависят друг от друга в реальном времени.

Kafka решает это иначе. Это распределённый лог сообщений: OrderService записывает событие «заказ создан», а все заинтересованные сервисы читают его самостоятельно — каждый в своём темпе. Если NotificationService временно упал, он просто догонит пропущенные события после перезапуска.

OrderService как producer пишет в Kafka, а из неё независимо читают три consumer-а: PaymentService, InventoryService и NotificationService

OrderService не знает, кто его читает. Новые потребители добавляются без изменений в коде отправителя.

Топик, партиция, offset

Топик

Сервисам нужно договориться, куда писать и откуда читать, не зная друг о друге. Точка встречи — топик: именованный лог сообщений, в который одни сервисы пишут, а другие из него читают. Имя отражает домен: order-events, payment-events, inventory-changes.

Каждое сообщение — пара (ключ, значение) плюс метаданные (timestamp, headers, offset):

ключ:    "order-12345"
значение: {"type":"OrderCreated","orderId":"12345","total":1500.00}

Партиция

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

Топик order-events разбит на три партиции, и в каждой лежит свой упорядоченный список сообщений от m0 и дальше

Зачем партиции:

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

Offset

Потребитель упал и поднялся — с какого места продолжать? Для этого каждое сообщение в партиции имеет порядковый номер — offset. Потребитель запоминает, до какого offset он дочитал, и продолжает именно оттуда. Это позволяет перечитать старые сообщения или начать с начала.

Из чего состоит кластер

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

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

Контроллер — роль, а не отдельный сервер: один из брокеров следит за составом кластера и распределением партиций, а при его отказе роль берёт другой. В старых версиях эту работу делал отдельный ZooKeeper, в современных (режим KRaft) всё живёт внутри самих брокеров, и ZooKeeper не нужен. Для читателя это значит, что в руководствах старше 2023 года будет лишний компонент, которого у вас нет.

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

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

Отправитель и потребитель

Producer — клиент, который пишет в Kafka. Consumer — клиент, который читает из Kafka.

// Отправитель
@Service
@RequiredArgsConstructor
public class OrderEventPublisher {
    private final KafkaTemplate<String, String> kafka;

    public void publish(OrderCreatedEvent event) {
        kafka.send("order-events", event.orderId(), toJson(event));
    }
}

// Потребитель
@Component
public class PaymentListener {
    @KafkaListener(topics = "order-events", groupId = "payment-service")
    public void onOrderCreated(ConsumerRecord<String, String> record) {
        OrderCreatedEvent event = fromJson(record.value(), OrderCreatedEvent.class);
        // обработка...
    }
}

Куда попадает сообщение: роль ключа

Producer выбирает партицию по правилу:

  1. Ключ указан → partition = hash(ключа) % количество_партиций. Все сообщения с одинаковым ключом всегда попадают в одну партицию.
  2. Ключ не указан → продьюсер сам решает, куда положить. Не по кругу, как можно подумать: он набирает пачку в одну партицию, отправляет её целиком и только потом переключается на следующую. Так получается меньше сетевых запросов, а нагрузка всё равно размазывается ровно.

Это важно для порядка. Если все события одного заказа должны быть в правильном порядке — используйте orderId как ключ:

// Правильно: ключ — orderId, все события заказа в одной партиции
kafka.send("order-events", event.orderId(), payload);

// Неправильно: без ключа события разлетятся по разным партициям
kafka.send("order-events", null, payload);

Число партиций не меняют на ходу

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

Порядок сообщений

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

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

Вывод: хотите упорядоченную обработку всех событий какого-то объекта — используйте его идентификатор как ключ сообщения.

Consumer Groups: как потребители делят работу

Consumer Group — это несколько потребителей с одинаковым groupId. Kafka распределяет партиции топика между ними: каждая партиция назначается ровно одному потребителю в группе.

Три партиции топика order-events, розданные группе payment-service: партиции 0 и 1 читает Consumer A, партицию 2 читает Consumer B

Что это даёт:

  • Параллелизм ограничен числом партиций: если партиций 3, а потребителей 5 — двое будут простаивать.
  • Каждое сообщение в группе обрабатывается ровно одним потребителем — нагрузка делится, а не дублируется.
  • Разные группы независимы: PaymentService и InventoryService обе читают order-events, но каждая со своим offset — как если бы у каждой была своя копия топика.

От топика order-events стрелки идут к трём независимым группам потребителей: payment-service, inventory-service и notification-service

Когда потребитель падает или добавляется новый, Kafka перераспределяет партиции — это называется rebalance. На время перераспределения группа не обрабатывает сообщения. Частые rebalance замедляют работу, поэтому стараются их избегать.

Ребаланс: почему группа периодически замирает

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

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

Живость проверяется двумя разными таймерами, и путать их не надо. session.timeout.ms (по умолчанию 45 секунд) — про сердцебиение: клиент отправляет служебные сигналы фоновым потоком, и если их нет столько времени, потребителя исключают. max.poll.interval.ms (по умолчанию 5 минут) — про обработку: если между двумя вызовами poll() прошло больше, группа считает, что потребитель завис, и забирает у него партиции. Второй и срабатывает в реальной жизни: пачка из пятисот сообщений, каждое по секунде — и вы вышли за пять минут, не сделав ничего плохого. Лечится это уменьшением max.poll.records (сколько сообщений отдаётся за один poll), а не увеличением таймаута до часа: большой таймаут означает, что настоящий зависший потребитель будет держать партиции час.

Два протокола распределения. Старый, eager: все потребители отпускают все партиции, потом получают новые назначения — на это время группа стоит целиком, и это то самое «замирание». Новый, cooperative-sticky: партиции переназначаются по частям, каждый потребитель отпускает только то, что действительно уходит другому, и остальные продолжают работать. В современных клиентах он и стоит по умолчанию (partition.assignment.strategy), а в старых его включают руками — это одна из самых дешёвых настроек, улучшающих поведение под нагрузкой.

Статическое членство группы. Отдельная беда — выкат: перезапустили три экземпляра по очереди, получили три ребаланса. Лечится настройкой group.instance.id: каждый экземпляр получает постоянное имя, и его короткий перезапуск (в пределах session.timeout.ms) не вызывает перераспределения — партиции ждут его возвращения. Для выкатов в Kubernetes, где имя пода стабильно у StatefulSet, это работает особенно удобно.

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

Commit offset: отметка о прогрессе

Потребитель сообщает Kafka, какие сообщения обработал — это называется commit offset. Без этого после перезапуска Kafka не знала бы, откуда продолжать.

Автоматический commit — режим «голого» клиента Kafka (enable.auto.commit=true). Работает он не по фоновому таймеру, как можно подумать, а внутри poll(): на каждом вызове клиент смотрит, прошло ли с прошлой фиксации auto.commit.interval.ms (по умолчанию 5 секунд), и если прошло — фиксирует всё, что успел отдать приложению. Удобно, но опасно: сообщения отданы, offset сдвинут, а обработаны они или нет — клиент не знает. Упали в этот момент — потеряли.

Если вы пишете на Spring Kafka, «по умолчанию» означает противоположное. Spring сам выставляет enable.auto.commit=false и фиксирует offset за вас — после того, как обработана вся пачка, выданная одним poll (режим AckMode.BATCH). Это уже не «раз в 5 секунд вслепую», но и не «после каждого сообщения».

Ручной commit — надёжнее. Только объект Acknowledgment прилетит в метод не сам по себе: контейнеру надо явно сказать, что подтверждает приложение, иначе параметр будет null и первый же ack.acknowledge() уронит обработчик.

spring:
  kafka:
    listener:
      ack-mode: manual   # без этой строки Acknowledgment в методе будет null
@KafkaListener(topics = "order-events", groupId = "payment-service")
public void onMessage(ConsumerRecord<String, String> record, Acknowledgment ack) {
    processEvent(record);
    ack.acknowledge(); // фиксируем только после успешной обработки
}

Try/catch здесь не нужен: если processEvent бросит исключение, до ack.acknowledge() дело не дойдёт, offset не сдвинется, и сообщение придёт снова.

auto.offset.reset: почему новый потребитель не видит старых сообщений

Первый вопрос, который задаёт каждый, кто только запустил потребителя: в топике лежат сообщения, потребитель подключился — и тишина. Дело не в сбое. Смещения у новой группы ещё нет, и потребитель должен решить, откуда читать; решает он это настройкой auto.offset.reset, и по умолчанию она равна latest — то есть «читать только то, что придёт дальше».

Значений два. latest начинает с конца: старое не читается вовсе. earliest начинает с самого начала топика: прочитано будет всё, что ещё не удалено по сроку хранения.

spring:
  kafka:
    consumer:
      auto-offset-reset: earliest

Работает эта настройка только когда смещения нет — у новой группы или когда сохранённое смещение уже вне доступного диапазона, потому что данные удалились по сроку хранения. Для группы, которая уже читала, она ничего не меняет: смещение есть, и потребитель продолжит с него. Отсюда типичная путаница: поставили earliest, перезапустили — и снова ничего, потому что группа та же и смещение сохранено. Чтобы перечитать заново, меняют имя группы или сдвигают смещение отдельной командой (kafka-consumer-groups --reset-offsets).

Выбор значения — это не вкус, а решение о смысле. Для аналитики и построения состояния из истории нужен earliest: пропустить события нельзя. Для обработки команд «здесь и сейчас» разумнее latest: после недельного простоя не нужно заваливать сервис недельным потоком. А вот earliest у второго случая даёт неприятный сюрприз в проде: сервис, поднятый с новой группой, начнёт заново обрабатывать всё, что хранится, и разошлёт письма за месяц.

С ручным commit'ом потребитель может получить одно сообщение дважды (при перезапуске после сбоя). Это называют at-least-once — и это стандартный подход. Цена: потребитель должен уметь обрабатывать одно и то же сообщение несколько раз без побочных эффектов (быть идемпотентным).

Заголовки и формат события

Сообщение в Kafka — это не только ключ и тело. У него есть заголовки: набор пар «имя — байты», как в HTTP. Тело при этом остаётся данными события, а всё служебное едет в заголовках.

var record = new ProducerRecord<>("order-events", order.id(), payload);
record.headers()
        .add("event-type", "OrderCreated".getBytes(UTF_8))
        .add("event-version", "2".getBytes(UTF_8))
        .add("traceparent", traceId.getBytes(UTF_8));

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

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

Дисциплины на это две. Проще — договорённость в команде плюс терпимый разбор JSON: неизвестные поля игнорируются, отсутствующие имеют значения по умолчанию. Строже — реестр схем (Schema Registry) с Avro или Protobuf, который сам отвергает несовместимое изменение при отправке; за это платят ещё одним сервисом в контуре. Выбор между ними — тот же разговор про цену эксплуатации, что и с самой Kafka; подробный разбор совместимости событий — в статье про шаблоны обмена сообщениями.

Гарантии доставки

Дубли и потери берутся с двух сторон. Отправитель не получил подтверждение и повторил отправку, хотя первая доехала, — дубль. Потребитель отметил прогресс раньше, чем обработал, и упал, — потеря; отметил позже и упал — дубль. Из того, что с этим делают, складываются три режима:

РежимСмыслПрименение
At-most-onceДоставлено или потеряно — без дублейМетрики, логи, телеметрия
At-least-onceДоставлено хотя бы один раз, возможны дублиБольшинство бизнес-сценариев
Exactly-onceРовно один раз, без потерь и дублейФинансовые операции, критичные к дублям

At-least-once + идемпотентный потребитель — стандартный выбор для бизнес-событий. Exactly-once технически реализуемо (transactional producer + read-committed consumer), но сложнее и медленнее.

Idempotent Producer

Producer может отправить сообщение дважды, если подтверждение не дошло: ответа нет, он повторяет отправку — а первая попытка на самом деле доехала. От этого защищает enable.idempotence=true: Kafka присваивает сообщениям порядковые номера и отбрасывает повторы. Начиная с Kafka 3.0 и эта настройка, и acks=all — значения по умолчанию, так что задача не «включить», а «не сломать»: выставленный вручную acks=1 «для скорости» заодно молча отключит и защиту от дублей.

spring:
  kafka:
    producer:
      properties:
        enable.idempotence: true   # с Kafka 3.0 это уже значение по умолчанию
        acks: all                  # и это тоже — здесь выписано явно, чтобы не потерять
        retries: 2147483647

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

Атомарность: запись в базу и отправка события

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

Решение — Transactional Outbox: событие записывается в ту же транзакцию базы данных, что и сам заказ. Отдельный фоновый процесс читает необработанные записи и отправляет их в Kafka.

OrderService одной транзакцией пишет заказ и запись outbox в PostgreSQL, фоновый @Scheduled polling-relay выбирает неопубликованные строки и шлёт их в Kafka идемпотентному потребителю

Ничего не теряется (запись атомарна), ничего лишнего не отправляется (relay видит только завершённые записи). Подробнее этот паттерн разобран в статье «Распределённые паттерны».

Репликация: что если сервер упадёт

Каждая партиция хранится не на одном сервере, а на нескольких. Один из них — лидер, остальные — реплики. Записи всегда идут к лидеру; читают обычно тоже с него, но потребителя можно настроить и на чтение с ближайшей реплики — это экономит трафик между зонами доступности, ценой небольшого отставания данных. Реплики догоняют лидера и формируют ISR (In-Sync Replicas — синхронизированные реплики). Синхронной считается та, что отстала от лидера не больше чем на replica.lag.time.max.ms — по умолчанию 30 секунд; отстала сильнее — брокер выбрасывает её из ISR, догнала — возвращает обратно. Если лидер падает, новым становится одна из реплик, оставшихся в ISR.

Producer указывает, насколько строгое подтверждение ему нужно (acks):

acksЧто значитРиск
0Отправил и забылМожно потерять
1Лидер получилПотеря при падении лидера до репликации
allВсе синхронизированные реплики получилиНадёжно

Для бизнес-событий — acks=all. Одного этого мало: настройка означает «дождись подтверждения от всех реплик, которые сейчас считаются синхронными», а если такая реплика осталась одна, то и подтверждение придёт от неё одной — и вместе с ней данные пропадут. Поэтому в паре с ней на топике задают min.insync.replicas=2: тогда при нехватке живых реплик запись честно завершится ошибкой, а не примет данные, которые некому сохранить.

acks=all в ISR один лидер лидер подтвердил лидер упал запись потеряна и min.insync=2 в ISR один лидер запись отклонена повтор отправки запись цела

Верхний ряд: при acks=all подтверждение приходит от единственной синхронной реплики и уходит вместе с ней; нижний: min.insync.replicas=2 отклоняет такую запись сразу.

Срок хранения сообщений

Kafka — это лог, и у него есть настраиваемый срок хранения:

  • По времени (retention.ms) — например, 7 дней.
  • По размеру (retention.bytes) — например, 1 ТБ на партицию.
  • Компакция (cleanup.policy=compact) — хранится только последнее сообщение для каждого ключа. Подходит для хранения текущего состояния: остатков на складе, актуальных цен, профилей пользователей. Удалить ключ из такого топика можно единственным способом — записать по нему сообщение со значением null. Его называют надгробием (tombstone): компакция сначала оставит его как последнее значение ключа, а спустя настроенное время уберёт и его. Без надгробий из compact-топика не исчезает ничего.

Это настройки самого топика, а не приложения: в application.yml их класть бесполезно, Kafka туда не смотрит. Задают их при создании топика или меняют у существующего:

# топик с компакцией — актуальные остатки
kafka-topics.sh --create --topic inventory-current \
  --bootstrap-server localhost:9092 \
  --config cleanup.policy=compact

# обычный топик событий — хранить 7 дней
kafka-topics.sh --create --topic order-events \
  --bootstrap-server localhost:9092 \
  --config cleanup.policy=delete --config retention.ms=604800000

# поменять срок хранения у уже существующего топика
kafka-configs.sh --alter --topic order-events \
  --bootstrap-server localhost:9092 \
  --add-config retention.ms=259200000

Когда Kafka не нужна

Kafka хорошо работает для потоков событий с высокой нагрузкой. Но она не для всего:

  • Нужен ответ на запрос — Kafka однонаправленная, запрос-ответ через неё громоздко. Лучше HTTP.
  • Мало событий — Kafka оптимизирована для высокой нагрузки; для нескольких событий в день это избыточно.
  • Нужна сложная маршрутизация по правилам — RabbitMQ с routing-правилами подойдёт лучше.

Список без замены бесполезен, поэтому вот чем заменяют, по случаям.

Фоновые задачи одного сервиса (отправить письмо, собрать отчёт, пересчитать кэш) — таблица заданий в той же базе, которую сервис и так использует: строка со статусом, выборка через SELECT ... FOR UPDATE SKIP LOCKED, повторы счётчиком попыток. Ни брокера, ни нового навыка эксплуатации, а транзакция общая с бизнес-данными. Это выбор по умолчанию для небольших сервисов, и его недооценивают.

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

Запрос с ответом — обычный HTTP или gRPC. Если смущает связанность, её лечат не брокером, а таймаутами, повторами и предохранителем; асинхронность ради асинхронности превращает простой вызов в распределённую отладку.

Оповещение браузера или мобильного приложения — WebSocket или отправка уведомлений через платформу, а не Kafka: клиенту не нужен журнал событий, ему нужно одно сообщение сейчас.

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

Дополнительно: при первом чтении можно пропустить

Глубже: сколько партиций и потребителей: расчёт, а не число с потолкарасширенное

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

Пропускная способность одного потребителя, C: сколько записей в секунду обрабатывает один экземпляр слушателя с настоящим обработчиком (запись в базу, вызов сервиса), измеряется на стенде. Целевой поток, T: пиковый, а не средний, с запасом на рост за срок жизни топика, обычно вдвое. Отсюда нижняя граница числа потребителей, T / C, а партиций не меньше, чем потребителей, иначе лишние стоят. Третья величина, целевое время разгребания: если после часового простоя потребителей накопилось B записей, а разгрести нужно за время R, потребителей нужно (B / R + T) / C, и это обычно больше, чем от потока. Четвёртая, поток одного продьюсера на партицию: лидер партиции принимает десятки мегабайт в секунду, и для потока в сотни мегабайт партиций нужно больше по этой причине, а не по потребителям.

Пример: обработчик тянет 200 записей в секунду, пик 1 000, час простоя даёт 3,6 миллиона записей, разгрести за полчаса: (3 600 000 / 1 800 + 1 000) / 200 = 15 потребителей, партиций 16 или 24 с запасом. Правило округления вверх до числа, которое делится на разумные размеры групп (6, 12, 24), чтобы партиции распределялись между потребителями поровну.

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

У RabbitMQ расчёт проще: очереди дёшевы, потребители масштабируются без потолка, и считают только T / C с запасом на пик, а надёжных очередей заводят десятки, не тысячи, потому что каждая это своя группа согласования.

Глубже: второй кластер и второй регион: MirrorMaker 2расширенное

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

Он работает как набор коннекторов Kafka Connect: читает топики из кластера-источника и пишет в кластер-приёмник, по умолчанию с префиксом имени источника (msk.orders вместо orders), чтобы при двусторонней репликации топики не зациклились; политику имён можно поменять на тождественную для активно-пассивной схемы. Порядок внутри партиции сохраняется, дубли возможны, потому что доставка между кластерами «хотя бы раз», и потребители на приёмнике должны быть идемпотентны так же, как в исходном кластере.

Главная сложность не в данных, а в оффсетах: у приёмника они свои, и группа потребителей, переехавшая на второй кластер, не знает, докуда дочитала на первом. MirrorMaker 2 ведёт служебный топик контрольных точек и переводит оффсеты групп между кластерами (с версии 2.7 может синхронизировать их сам), так что после переключения группа продолжает с примерно того же места, с повторами вместо потерь. Отставание репликации между кластерами это и есть потеря при аварии, RPO, и его выносят на график.

Две схемы. Активно-пассивная: пишут в один кластер, второй получает копию и ждёт; переключение это смена адресов у продьюсеров и потребителей плюс перевод оффсетов, и его репетируют заранее, потому что вручную под давлением оно не получается. Активно-активная: каждый регион пишет в свой кластер, а читает и свои, и зеркальные топики соседа, объединяя их в приложении; дубли и порядок между регионами решает сама бизнес-логика, и это на порядок сложнее. У облачных и коммерческих дистрибутивов есть свои средства связи кластеров без префиксов и с сохранением оффсетов, но идея та же: асинхронная копия с отставанием и переключение как процедура. У RabbitMQ роль MirrorMaker играют Federation и Shovel, о чём статья про RabbitMQ в production.

Коротко

  • Топик — именованный лог событий. Партиция — его физическая часть, упорядоченная и неизменяемая. Ключ сообщения определяет партицию: один ключ → одна партиция → порядок гарантирован.
  • Offset — позиция потребителя. Commit offset после обработки, не до. Consumer Group делит партиции между потребителями. Несколько групп читают один топик независимо.
  • Стандартный выбор: at-least-once + ручной commit + идемпотентный потребитель. Для атомарности «база + Kafka» — паттерн Transactional Outbox.
  • Репликация защищает от потери данных при падении сервера. Для надёжности — acks=all вместе с min.insync.replicas=2. Срок хранения настраивается по времени или размеру; компакция хранит только последнюю версию для каждого ключа.
  • Партиции считают от измеренного: (накопленное / время разгребания + пиковый поток) / поток одного потребителя, с запасом и кратно размеру группы; сверху ограничивают файлы и восстановление брокера, а горячий ключ не лечится числом партиций.
  • Между кластерами и регионами данные возит MirrorMaker 2: асинхронно, с префиксом топиков, «хотя бы раз», с переводом оффсетов групп; отставание это RPO, переключение репетируют.
  • Внутри кластера три роли: брокеры хранят партиции, один из них контроллер (в режиме KRaft, без ZooKeeper), а запись и чтение идут через лидера партиции; bootstrap.servers — только точка знакомства, дальше клиент ходит к брокерам напрямую.
  • Новая группа по умолчанию читает только новое (auto.offset.reset: latest), а earliest действует лишь когда смещения нет — перечитать заново можно сменой группы или сбросом смещений.
  • Ребаланс чаще всего вызывает долгая обработка (max.poll.interval.ms), лечится он меньшим max.poll.records, протоколом cooperative-sticky и статическим членством через group.instance.id.
  • Служебное едет в заголовках (тип, версия, трассировка), а формат тела меняют только совместимо: новые поля необязательные, удаление в два шага, дисциплину держит реестр схем или терпимый разбор.

Что пощупать

Ключ, заголовки, порядок и атомарность записи с отправкой в практикуме remodov/marketplace-system собраны в одном коротком классе: сервис заказов публикует события из таблицы исходящих в один топик, ключом ставит идентификатор заказа, а метаданные события кладёт в заголовки. Kafka поднимается из infra/compose.yaml.

Код: adapter-out-kafka, контракт.

Сделаем сами

Ветка step-10-events-and-contract — публикация в outbox и правила сериализации вынуты, три теста красные.

Что почитать дальше

  • Kafka в production — Spring Kafka, Dead Letter Queue, Schema Registry, consumer lag, настройка производительности, безопасность, KRaft.
  • Распределённые паттерны — Outbox + polling-relay, Saga, Idempotent Consumer подробно.
  • CQRS — Kafka как канал между командной и читающей моделью.