Distributed Message Queue (Kafka-style) case study. Producers (idempotent + transactional) write batched/compressed records to a 3-broker cluster (DC-1) with RF=3 and min.insync.replicas=2. Three partitions (P0/P1/P2) each have a leader and 2 followers spread across racks. Control plane is KRaft metadata quorum (no ZooKeeper) with 3 voters, plus group coordinator and transaction coordinator. Two consumer groups: analytics (auto.commit=false, 1:1 partition assignment) and alerting (read_committed for EOS). Internal compacted topics: __consumer_offsets, __transaction_state, __cluster_metadata. Multi-DC mirror via MirrorMaker 2 to DC-2, with tiered storage (S3) for cold segments. Two ADR panels: pull vs push consumer model, KRaft vs ZooKeeper metadata. Six animated scenarios: produce + ISR replicate (acks=all), consumer group fetch + offset commit (zero-copy sendfile), leader broker failure + ISR election (KRaft fast failover), replay from offset (retention + tiered recovery), exactly-once via transactional producer (PID/epoch fence + 2PC commit markers + read_committed), and ADR walkthrough.
Распределенная очередь сообщений нужна, когда система должна принимать события быстрее, чем downstream-сервисы успевают их обработать, и при этом не терять данные. Это базовый компонент event-driven architecture: платежи публикуют события для аналитики, логирование пишет миллионы записей в секунду, CDC переносит изменения из базы, ML pipeline читает историю повторно, а consumer groups независимо обрабатывают один и тот же поток.
Kafka-like очередь важна тем, что это не просто брокер с push-доставкой. Это распределенный commit log с partition ordering, durable retention, replay by offset, consumer groups, replication и контролируемой моделью подтверждений. На интервью этот кейс проверяет понимание диска, page cache, batching, backpressure, consumer lag, leader election, quorum writes и exactly-once semantics.
Масштаб кейса: десятки тысяч topics, сотни тысяч partitions, сотня brokers, миллионы сообщений в секунду, replication factor 3, retention на дни и многократный fan-out на consumer groups. При таком масштабе bottleneck часто не CPU, а сеть, disk bandwidth, количество partitions на broker, размер batches и количество rebalance events.
Связанные темы: ::concept{slug="message-queues"}, ::concept{slug="kafka-deep"}, ::concept{slug="replication"}, ::concept{slug="consensus-overview"}, ::concept{slug="partitioning-strategies"}.
Ментальная модель Kafka-like системы: topic разбит на partitions, каждая partition является append-only log. Producer выбирает partition по key hash или custom partitioner, отправляет batch лидеру partition, лидер пишет batch в segment file, followers догоняют лидер через replication fetch, а consumer читает по offset.
Ordering гарантируется только внутри одной partition. Если все сообщения одного заказа имеют одинаковый key order_id, они попадут в одну partition и сохранят порядок. Если key не задан или распределяется случайно, глобального порядка нет. Это не bug, а цена горизонтального масштабирования.
Consumer group превращает partitions в единицы параллелизма. Внутри одной группы partition назначается одному consumer instance, чтобы сохранить порядок. Несколько групп могут читать один и тот же topic независимо: analytics, alerting, billing reconciliation. Offset хранится отдельно от данных, поэтому consumer может replay старый диапазон, если retention еще не удалил сегменты.
Durability задается связкой acks=all, replication factor и min.insync.replicas. Если RF=3 и min ISR=2, запись считается успешной, когда лидер и хотя бы один follower подтвердили запись. Это переживает потерю одного broker, но не гарантирует выживание при потере двух реплик до восстановления.
Диаграмма показывает producers, broker cluster, KRaft control plane, consumer groups, internal compacted topics, tiered storage и multi-DC mirror. Producers не пишут во все brokers подряд: они получают metadata, узнают лидера partition и отправляют batch именно туда. Brokers держат лидеров и followers для разных partitions, чтобы нагрузка распределялась по кластеру.
Control plane через KRaft хранит metadata: кто leader, кто в ISR, какие topics и partitions существуют, какие quotas и configs активны. Это не data plane; поток сообщений идет через brokers, а controller управляет topology и elections. Отдельные coordinators обслуживают consumer groups и transactions.
Internal topics важны для понимания системы. __consumer_offsets хранит committed offsets, __transaction_state хранит состояние transactional producers, metadata topic хранит cluster state. Это обычные log-based structures, но compacted, потому что нужна последняя версия ключа.
Tiered storage показывает, что старые segments можно выгружать в S3-like хранилище. Это удешевляет долгий retention, но делает replay старой истории медленнее.
Produce + ISR replicate учит, что запись сначала идет лидеру partition. Лидер appends batch в log, followers подтягивают данные, ISR обновляется, и только после нужного числа подтверждений producer получает ACK. Batching и compression резко повышают throughput, но добавляют latency через linger time.
Consumer fetch показывает pull model. Consumer сам запрашивает диапазон offsets и контролирует скорость. Broker может отдать данные через page cache и zero-copy sendfile, поэтому последовательное чтение из log очень эффективно. Offset commit отделен от чтения: consumer может обработать batch и только потом зафиксировать progress.
Leader failure показывает выбор нового лидера из ISR. Если лидер умер, controller выбирает follower, который был синхронен. Producer получает metadata error, обновляет metadata и повторяет запрос. Если ISR пуст или выбран unclean leader, появляется риск потери подтвержденных сообщений.
Replay from offset учит, что очередь не удаляет сообщение после чтения. Retention управляется временем и размером, а consumer position независим. Это делает Kafka удобной для reprocessing, backfill и восстановления downstream materialized views.
Exactly-once scenario показывает границы EOS. Idempotent producer устраняет duplicate при retries, transactional producer атомарно пишет в несколько partitions/topics и commit markers, consumer с read_committed не видит aborted records. Но внешняя БД все равно требует idempotency или transactional outbox.
Pull против push. Pull дает consumer backpressure и простой batching, но latency зависит от poll interval и fetch settings. Push удобен для простых queues, но сложнее контролировать перегрузку клиентов.
Partitions дают параллелизм, но слишком много partitions ухудшает metadata, recovery, memory usage и rebalance time. Недостаточно partitions ограничивает throughput и количество consumers в группе.
acks=1 быстрее, но может потерять данные при падении лидера до replication. acks=all с min ISR надежнее, но медленнее и может отвергать запись, когда кластер деградировал.
Long retention повышает replay capability, но требует storage и compaction/tiered storage. Compacted topics сохраняют последнюю версию ключа, но не подходят для полного audit log без отдельной retention policy.
Exactly-once повышает корректность pipeline, но добавляет coordination, transaction state, fencing и operational complexity. Для многих analytics задач at-least-once плюс idempotent consumer дешевле и понятнее.
Apache Kafka является главным примером: partitions, brokers, ISR, consumer groups, KRaft, log segments, offset index, transactions, compacted topics. Redpanda реализует Kafka API на C++/Seastar и делает ставку на меньшую operational complexity. Apache Pulsar разделяет brokers и BookKeeper storage, имеет segment storage и multi-tenancy. RabbitMQ лучше подходит для routing, work queues и lower-throughput task dispatch, но не заменяет Kafka для replayable event log.
Amazon Kinesis похож концептуально: shards, sequence numbers, retention, consumers. Google Pub/Sub ближе к managed messaging с другой моделью delivery и ack. Для CDC часто используется Debezium -> Kafka -> consumers.
Первый anti-pattern: считать Kafka обычной job queue и создавать один topic на каждую задачу без модели retention, partition key и consumer group ownership.
Второй anti-pattern: использовать случайный partition key для сущностей, которым нужен порядок. Если события заказа попадают в разные partitions, consumer увидит их в произвольном порядке.
Третий anti-pattern: commit offset до обработки. При crash сообщение будет считаться обработанным, хотя side effect не выполнен. Лучше commit после обработки или использовать transactional pattern.
Четвертый anti-pattern: бесконтрольно увеличивать partitions. Это может ухудшить failover, controller load и rebalance, хотя на бумаге увеличивает параллелизм.
Пятый anti-pattern: обещать exactly-once end-to-end только потому, что broker поддерживает transactions. Внешние API, email, платежи и SQL updates требуют идемпотентности, outbox или deduplication keys.
Не используйте Kafka-like очередь для синхронного request/response, где клиент ждет ответ за 50 мс. HTTP/gRPC проще и понятнее.
Не используйте distributed log для маленького монолита с десятками сообщений в минуту, если достаточно Postgres table queue или managed task queue. Operational cost Kafka может быть выше пользы.
Не используйте Kafka как primary database для произвольных point reads. Log оптимизирован для последовательного append/read и replay, а не для произвольных запросов по secondary indexes.
Не используйте ее для задач, где сообщение должно быть удалено сразу после одного consumer. Для классического work queue с конкурирующими workers RabbitMQ/SQS/Celery-подобная модель может быть проще.
Разберите Kafka protocol basics, producer batching, compression, partitioner, ISR, leader election, KRaft metadata quorum, consumer group protocol и offset commits. Затем изучите log compaction, tiered storage, idempotent producer, transactions, read_committed, MirrorMaker 2 и schema registry. Для архитектурных связок полезны ::concept{slug="change-data-capture"}, ::concept{slug="event-driven-architecture"}, ::concept{slug="replication"} и ::concept{slug="consensus-overview"}.