Created by AIImproved by people

Riqli · living documents · updated continuously

Where AI knowledge meets human practice.

Share with friends
image

Working with message queues (Kafka, RabbitMQ, ActiveMQ). Hard skill. (Message queues. Kafka. RabbitMQ. ActiveMQ. Publish-subscribe. Producer-consumer. Event-driven. Self-study. Q&A. Tutorials. Documentation.)

sections

Введение

contents

Аннотация. В современном мире распределённых систем и микросервисной архитектуры вопрос надёжного и эффективного обмена данными между компонентами становится критическим. Синхронные HTTP-запросы часто приводят к каскадным отказам, высокой задержке и жёсткой связанности сервисов. На смену им приходят брокеры сообщений — промежуточное программное обеспечение, обеспечивающее асинхронное взаимодействие. Данный курс посвящён троице лидеров в этой области: Apache Kafka, RabbitMQ и ActiveMQ. Мы не просто перечислим их возможности, а разберём архитектурные принципы, сценарии применения, показатели производительности и критерии выбора. Этот материал будет полезен как системным аналитикам, формирующим требования к интеграциям, так и разработчикам, которые хотят сделать осознанный технологический выбор и избежать типичных ошибок эксплуатации. Курс построен на основе практических кейсов, сравнительных бенчмарков и современных тенденций, таких как гибридные подходы и облачные развёртывания .

contents

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

contents

Результаты обучения.

  • Знать: Архитектурные различия моделей «очередь» и «лог»; принципы работы AMQP и KRaft; назначение компонентов (Exchange, Partition, Consumer Group).
  • Уметь: Проектировать маршрутизацию сообщений с помощью обменников RabbitMQ; настраивать группы потребителей и партиционирование в Kafka; интегрировать ActiveMQ в Java-приложения на основе JMS.
  • Владеть: Навыками интерпретации графиков задержек и пропускной способности; пониманием стратегий обработки сбоев (Dead Letter Queue, Retry); готовностью к реализации паттернов событийного проектирования.
contents

Для кого этот курс.

Курс предназначен для системных архитекторов, разработчиков бэкенда (особенно на Java, Python и Go), а также системных аналитиков, участвующих в проектировании распределённых систем. Если вы сталкивались с падением микросервисов из-за перегрузки или необходимостью обработать миллионы событий от IoT-устройств — этот курс даст вам язык и инструменты для решения проблемы .

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

sections

Модуль 1. АРХИТЕКТУРНЫЕ ПАРАДИГМЫ: ОЧЕРЕДЬ ПРОТИВ ЛОГА

contents

Эволюция брокеров: от JMS к событийным потокам. Первое поколение систем (IBM MQ, TIBCO) решало задачи интеграции корпоративных приложений, но страдало от низкой производительности и сложности кластеризации. Спецификация JMS (Java Message Service) стандартизировала работу с очередями в Java-мире, что породило такие реализации, как ActiveMQ. Второе поколение, представленное RabbitMQ и протоколом AMQP 0-9-1, внедрило концепцию «умного брокера», способного гибко маршрутизировать сообщения. Третье поколение — Apache Kafka — переосмыслило брокера как распределённый лог-систему (commit log), отказавшись от удаления сообщений после доставки в пользу долгосрочного хранения. Это разделение на «очереди» (с удалением) и «логи» (с хранением и воспроизведением) является фундаментальным при выборе архитектуры .

contents

Модель очереди (Queue) на примере RabbitMQ и SQS. В классической модели очереди сообщение существует до момента его подтверждения (ACK) единственным потребителем. Если потребитель подтвердил получение, сообщение удаляется навсегда. Если потребитель упал до отправки ACK, сообщение возвращается в очередь для повторной обработки (режим «видимости» в SQS). RabbitMQ реализует это через очереди (Quorum Queues на Raft для отказоустойчивости), а ActiveMQ — через JMS-очереди. Основная метафора здесь — конкурирующие потребители (Competing Consumers), которые позволяют распределять нагрузку по задачам (Job Queues) . Главное ограничение: если через месяц появится новый сервис, он не сможет «перечитать» историю сообщений, так как они уже удалены.

contents

Модель лога (Log) на примере Kafka. Apache Kafka представляет данные как неизменяемый, упорядоченный журнал записей. Сообщения внутри топика распределяются по партициям (Partitions). Каждое сообщение имеет смещение (offset). Потребители не удаляют сообщения; они просто сдвигают указатель смещения в своей группе. Благодаря этому новые потребители могут начать читать с самого начала или с любой точки в прошлом (replay). Эта модель идеальна для CDC (Change Data Capture), источников истины и аналитики. Kafka не маршрутизирует сообщения сложным образом (как RabbitMQ), а просто сохраняет их в партициях, полагаясь на потребителей («умный потребитель» против «умного брокера») .

contents

Сравнение гарантий доставки (At-most-once, At-least-once, Exactly-once). Это один из ключевых критериев выбора. At-most-once (огне-и-забудь) — сообщение может потеряться, но дублей не будет (высокая скорость). At-least-once — стандарт в брокерах: производитель ждёт подтверждения (ACK), но при сбое брокера или сети сообщение может быть доставлено дважды. Потребитель должен быть идемпотентным. Exactly-once достижим только в ограниченных условиях (например, в транзакциях Kafka между Kafka-топиками) и требует включения enable.idempotence=true. В RabbitMQ точно-один-раз технически невозможен без дополнительных механизмов координации, но подтверждения (publisher confirms + consumer acks) дают уверенность в at-least-once .

contents

Физическая модель данных: хранилище и индексы. В RabbitMQ данные хранятся в очередях, которые представляют собой структуры в памяти и на диске. При использовании Quorum Queues данные реплицируются между узлами (Raft). В ActiveMQ хранение по умолчанию — KahaDB (файловое хранилище). В Kafka сообщения записываются на диск последовательно (sequential I/O), что является секретом её высокой пропускной способности. Каждая партиция — это отдельный файл-сегмент на файловой системе. Запись в конец файла (append) позволяет избежать дорогостоящих операций случайного ввода-вывода, характерных для очередей, что дает выигрыш в скорости при больших объёмах .

sections

Модуль 2. ГЛУБОКОЕ ПОГРУЖЕНИЕ В APACHE KAFKA

contents

Основные компоненты: Producer, Topic, Partition, Consumer Group. Producer публикует данные в Topic. Топик делится на Partitions для параллелизма. Порядок сообщений гарантируется строго ВНУТРИ одной партиции, но не между разными партициями. Consumer Group — это набор потребителей, которые совместно читают топик. Каждая партиция в группе обслуживается только одним потребителем. Если потребителей больше, чем партиций, часть простаивает. Это ключевое ограничение масштабируемости: максимальный уровень параллелизма ограничен количеством партиций .

contents

KRaft: отказ от ZooKeeper (Kafka Raft). До версии 3.x Kafka полагалась на внешнюю систему координации ZooKeeper, что усложняло развёртывание. Начиная с версии 4.0 (и в современных релизах), ZooKeeper полностью удалён. Вместо него используется встроенный сервис KRaft (Kafka Raft). Теперь контроллеры (контроллеры) сами управляют метаданными кластера (кто является лидером партиции) через консенсус Raft. Это упрощает настройку, снижает задержки при перевыборах лидера и уменьшает количество движущихся частей в инфраструктуре .

contents

Стратегии партиционирования и ключи сообщений. При отправке сообщения вы можете указать ключ (key). Хеш ключа определяет, в какую партицию попадёт сообщение. Это гарантирует, что все сообщения с одним ключом (например, user_id) попадут в одну партицию и будут прочитаны в том порядке, в котором были отправлены. Если ключ не указан, используется стратегия round-robin или sticky-партиционирование (пакетная отправка в одну партицию для уменьшения накладных расходов). Неправильный выбор ключа может привести к перекосу данных (skew), когда одна партиция переполнена, а другие пусты .

contents

Асинхронная отправка и обратные вызовы. Производитель Kafka поддерживает асинхронную отправку через колбэки (send(record, callback)). Это позволяет отправлять миллионы сообщений в секунду, не блокируя основной поток на ожидание подтверждения. При этом важно обрабатывать колбэки для логирования ошибок. Частая ошибка — использование синхронного метода send().get(), который убивает производительность. Для повышения пропускной способности также используются параметры batch.size и linger.ms, накапливающие сообщения перед отправкой .

contents

Управление смещениями (Consumer Offsets). Потребители коммитят (сохраняют) смещение прочитанных сообщений. Если обработка сообщения завершилась ошибкой, можно либо закоммитить смещение и перейти к следующему (пропуск), либо не коммитить — тогда сообщение будет прочитано снова (при перебалансировке). Стратегия auto.offset.reset определяет поведение при отсутствии закоммиченного смещения: earliest (читать всё сначала) или latest (читать только новые). Ручное управление смещениями (manual commit) предпочтительнее для критичных данных, чтобы обеспечить семантику «обработано не более одного раза» в связке с идемпотентностью .

contents

Тир-уровневое хранение и долгосрочная ретенция (Tiered Storage). В новых версиях Kafka появилась поддержка tiered storage. Это означает, что старые данные можно автоматически выгружать из дорогого кластера брокеров в дешёвое облачное хранилище (например, S3), сохраняя при этом возможность их чтения. Это позволяет держать ретенцию неделями или месяцами без перегрузки дисков брокеров, что критично для аудита и повторного анализа исторических данных .

contents

Стриминг-обработка: Kafka Streams и ksqlDB. Kafka — это не только очередь, но и платформа для потоковой обработки. Kafka Streams — это Java-библиотека для построения приложений реального времени (подсчёты статистики, обогащение потоков) без необходимости в отдельных кластерах обработки (в отличие от Flink). ksqlDB позволяет выполнять SQL-подобные запросы к потокам. Это делает Kafka центром обработки данных, а не просто трубопроводом, что отличает её от RabbitMQ или ActiveMQ, которые фокусируются только на доставке .

sections

Модуль 3. ГЛУБОКОЕ ПОГРУЖЕНИЕ В RABBITMQ

contents

Концепция Exchanges, Bindings и Routing Keys. В отличие от Kafka, где сообщение отправляется прямо в топик-лог, в RabbitMQ производитель отправляет сообщение в Exchange (обменник). Это уровень маршрутизации. Binding — это правило, связывающее обменник и очередь через Routing Key. Обменник получает сообщение, смотрит на его routing key и, сверяясь с биндингами, решает, в какие очереди его положить. Это позволяет реализовать сложную логику маршрутизации (например, направить все сообщения с префиксом «order» в одну очередь, а с «payment» — в другую) .

contents

Типы обменников: Direct, Topic, Fanout, Headers. Direct Exchange — точное совпадение ключа маршрутизации (как почтовый адрес). Topic Exchange — маршрутизация по шаблону с использованием масок (* и #), например, *.error для всех ошибок. Fanout Exchange — широковещательная рассылка во все привязанные очереди (игнорирует ключ), используется для обновления кэшей или уведомлений всех сервисов. Headers Exchange — маршрутизация на основе заголовков сообщений (а не ключа), что удобно, когда ключи сложно структурировать .

contents

Надёжность и подтверждения (Publisher Confirms и Consumer Acks). RabbitMQ гарантирует надёжность через двустороннее рукопожатие. Производитель включает режим publisher confirms, и брокер присылает подтверждение (ACK) только после того, как сообщение записано на диск и продублировано (если настроена зеркальная очередь). Потребитель должен отправить ACK брокеру после успешной обработки. Если потребитель упал без ACK, сообщение возвращается в очередь. Это даёт огромную гибкость, но требует, чтобы потребители были идемпотентными в случае повторной доставки .

contents

Quorum Queues: отказ от зеркалирования. Ранее в RabbitMQ для высокой доступности использовались Mirrored Queues. Они были сложны в управлении и имели проблемы со «split brain». Начиная с версии 3.8, стандартом стали Quorum Queues, построенные на основе консенсуса Raft. Они обеспечивают более безопасную репликацию данных между узлами, автоматическое восстановление и лучшее поведение при сетевых сбоях. Все современные инсталляции RabbitMQ должны использовать именно Quorum Queues для критичных данных .

contents

Dead Letter Exchanges (DLX) и задержки. Если сообщение не может быть доставлено потребителю (истек TTL, очередь переполнена, отклонено с reject), оно может быть отправлено в Dead Letter Exchange. Это позволяет организовать обработку ошибок: например, после трёх неудачных попыток перенаправлять сообщение в очередь для ручного анализа. Для реализации отложенных задач (задержка в 5 минут) часто используется плагин rabbitmq_delayed_message_exchange, так как нативные TTL требуют создания очереди-сброса, что менее гибко .

contents

Prefetch и управление потоком (Flow Control). Параметр prefetch определяет, сколько сообщений брокер может «толкнуть» потребителю без подтверждения. Если установить prefetch=1, потребитель будет получать по одному сообщению за раз (справедливое распределение, нет перегрузки медленного воркера). Если установить высокое значение (например, 100), производительность вырастет, но если один потребитель упадёт, эти 100 сообщений будут отправлены повторно. RabbitMQ также использует механизм back pressure (Credit-based flow control), замедляя производителей, если очереди переполнены или потребители не успевают .

sections

Модуль 4. ACTIVE MQ: КЛАССИЧЕСКИЙ JMS-БРОКЕР

contents

Место ActiveMQ в экосистеме и архитектура. ActiveMQ — это «ветеран» брокеров сообщений, написанный на Java. Его главная сила — строгое следование спецификации JMS (Java Message Service), что делает его идеальным выбором для enterprise-систем, стандартизированных на Java EE. Архитектурно он представляет собой классический брокер с очередями и топиками. Он поддерживает множество протоколов: AMQP, MQTT, STOMP, OpenWire. Однако его производительность (десятки тысяч сообщений в секунду) значительно уступает Kafka и RabbitMQ, поэтому он постепенно вытесняется в высоконагруженных сценариях .

contents

Персистентность в ActiveMQ: KahaDB и JDBC. По умолчанию ActiveMQ использует хранилище KahaDB — файловый журнал, оптимизированный для быстрой записи и восстановления. Это быстрее, чем запись в реляционную базу данных. Тем не менее, ActiveMQ поддерживает сохранение сообщений в JDBC-хранилища (MySQL, PostgreSQL) для тех, кому важна централизованная семантика транзакций. Важно знать, что JDBC-хранилище значительно медленнее и создает узкое место на уровне базы данных, поэтому оно рекомендуется только для низких нагрузок .

contents

JMS vs AMQP: выбор протокола. Хотя ActiveMQ и поддерживает AMQP, его «родным» языком является JMS. JMS — это API-спецификация (интерфейс) для Java, а не сетевой протокол, как AMQP. Это означает, что взаимодействие между Java-приложениями через JMS происходит эффективно (используя внутренний протокол OpenWire), но интеграция с приложениями на Python или C# (через AMQP 1.0) может быть менее оптимальной или иметь отличия в семантике. В то время как RabbitMQ изначально построен на AMQP и обеспечивает бесшовную работу с разными языками, ActiveMQ ориентирован на Java-стек .

contents

Сценарии применения ActiveMQ сегодня. В 2026 году ActiveMQ редко встречается в «зелёных полях» (Greenfield) проектах. Его ниша — это устаревшие (Legacy) системы финансового сектора и госсектора, где жестко регламентировано использование JMS. Также он используется как легковесное решение для тестовых стендов или небольших проектов, где не требуется миллионная нагрузка. Однако его сообщество менее активно, чем у RabbitMQ или Kafka, и новые фичи появляются редко. Для новых проектов рекомендуется рассматривать ActiveMQ Artemis — преемника с улучшенной производительностью, либо сразу RabbitMQ/Kafka .

sections

Модуль 5. СРАВНИТЕЛЬНЫЙ АНАЛИЗ И БЕНЧМАРКИНГ

contents

Пропускная способность (Throughput) и задержка (Latency). Исследования показывают принципиальную разницу . Kafka демонстрирует феноменальную пропускную способность — до 1.2 миллионов сообщений в секунду, с задержкой p95 около 18 мс. RabbitMQ показывает задержки менее 1 мс при низкой нагрузке, но пропускная способность падает до 10-50 тысяч сообщений в секунду при росте нагрузки (особенно при персистентности). ActiveMQ обычно отстаёт от RabbitMQ по скорости. Причина: Kafka пишет в последовательный лог (быстро), а RabbitMQ — выполняет сложные маршрутизационные операции и операции с индексом очереди .

contents

Влияние персистентности на производительность. Включение долговременного хранения (persistence) значительно влияет на скорость. В Kafka запись на диск является обязательной и является её сильной стороной (optimized for disk). В RabbitMQ персистентные сообщения (с флагом delivery_mode=2) пишутся на диск синхронно (если не настроен lazy queue), что может вызвать падение пропускной способности до 5-10 раз по сравнению с режимом транзитной очереди. ActiveMQ также сильно проседает в производительности при включении персистентности через JDBC .

contents

Масштабируемость: вертикальная vs горизонтальная. Kafka масштабируется горизонтально, добавляя брокеры. Производительность растёт линейно с добавлением партиций и узлов. RabbitMQ также поддерживает кластеризацию, но масштабирование сложнее из-за необходимости репликации состояния очередей (Quorum Queues). В нем часто требуется вертикальное масштабирование (увеличение RAM/CPU) для одной очереди, так как одна очередь не может быть распределена по нескольким узлам (она принадлежит одному мастеру). ActiveMQ поддерживает master-slave кластеризацию, но её горизонтальные возможности ограничены .

contents

Операционная сложность (Ops Complexity). RabbitMQ и ActiveMQ считаются проще в установке и эксплуатации (особенно через Amazon MQ). Kafka, несмотря на KRaft, по-прежнему считается сложной системой: нужно мониторить репликацию, следить за дисковым пространством (так как хранит данные долго), управлять партициями и знать внутренние параметры JVM. Согласно данным Stack Overflow, 10.9% разработчиков предпочитают RabbitMQ из-за его понятности .

contents

Сравнительная таблица характеристик (Kafka, RabbitMQ, ActiveMQ).

  • Модель: Kafka (Распределённый лог), RabbitMQ (Очередь/Обменник), ActiveMQ (Очередь/Topic JMS).
  • Протокол: Kafka (TCP binary), RabbitMQ (AMQP 0-9-1 / 1.0, MQTT), ActiveMQ (OpenWire, JMS, AMQP, MQTT).
  • Порядок: Kafka (Внутри партиции), RabbitMQ (Внутри очереди), ActiveMQ (Внутри очереди).
  • Ретенция: Kafka (Дни/месяцы или бесконечно), RabbitMQ (До подтверждения потребителем), ActiveMQ (До подтверждения).
  • Воспроизведение: Kafka (Да), RabbitMQ (Нет, после ACK удалено), ActiveMQ (Нет).
  • Высокая доступность: Kafka (Репликация партиций), RabbitMQ (Quorum Queues / Raft), ActiveMQ (Master-Slave).
sections

Модуль 6. ПРАКТИЧЕСКАЯ ИНЖЕНЕРИЯ И СЦЕНАРИИ

contents

Паттерн «Конкурирующие потребители» (Competing Consumers). Это базовый паттерн для увеличения пропускной способности. В RabbitMQ несколько инстансов одного сервиса читают из одной очереди. Сообщения распределяются между ними (обычно round-robin). Это позволяет горизонтально масштабировать обработчики, но важно следить за prefetch, чтобы медленный потребитель не накапливал у себя огромный батч. В Kafka аналогия — это несколько потребителей в одной Consumer Group, где каждый читает свою партицию. В отличие от RabbitMQ, здесь нет состязания за одно сообщение; оно гарантированно попадает в одного потребителя (по partition) .

contents

Паттерн «Очередь мёртвых писем» (Dead Letter Queue - DLQ). Все три брокера поддерживают DLQ, но по-разному. В RabbitMQ это реализуется через Dead Letter Exchange, куда попадают сообщения, отвергнутые (basic.nack) или с истекшим TTL. В Kafka DLQ часто реализуется как отдельный топик topic-name-dlq, куда потребитель пишет сообщения после нескольких неудачных попыток (retries). В ActiveMQ это встроенная функция JMS (просроченные сообщения). DLQ критичны для диагностики ошибок и предотвращения блокировки основной очереди «плохими» сообщениями (poison pills) .

contents

Гибридный подход: Kafka + RabbitMQ в одной системе. Это современный тренд . Например, RabbitMQ используется для обработки транзакционных заказов (сложная маршрутизация, гарантии доставки), а Kafka — для сбора логов и аналитики на основе этих же событий. Данные из RabbitMQ могут быть стримлены в Kafka через плагины (например, rabbitmq-kafka-bridge) для долгосрочного хранения и анализа. Это позволяет использовать сильные стороны каждого инструмента: гибкость RabbitMQ и высокую пропускную способность Kafka.

contents

Мониторинг и наблюдаемость (Observability). Ключевые метрики для мониторинга: размер очереди (Messages Ready), количество неподтверждённых (Unacked), наличие потребителей (Consumer Count) . В RabbitMQ важно следить за памятью (Memory Alarm) и диском (Disk Alarm). При превышении лимитов брокер блокирует продюсеров. В Kafka — следить за смещением потребителей (consumer lag), чтобы понять, что потребитель отстаёт от производителя. Grafana и облачные дашборды (Amazon MQ) предоставляют готовые решения для визуализации .

contents

Cloud-native и управляемые сервисы: Amazon MQ, Confluent Cloud. В AWS существует сервис Amazon MQ, который управляет брокерами RabbitMQ и ActiveMQ. Это снимает бремя администрирования кластеров. Для Kafka популярны Confluent Cloud и MSK (Managed Streaming for Kafka). В облачных решениях важно учитывать стоимость: за запросы в SQS (не рассматривается, но аналогично) и за хранение в Kafka. Использование управляемых сервисов сокращает операционные расходы (Ops), но увеличивает стоимость вендор-локина .

contents

Критерии выбора и дерево принятия решений.

  • Если нужна сложная маршрутизация (topics, headers) и RPCRabbitMQ.
  • Если нужна высокая пропускная способность (100k+ msg/s), долгосрочное хранение и replayKafka.
  • Если проект строго на Java EE с JMS, низкая нагрузкаActiveMQ.
  • Если вы в AWS и хотите простой менеджмент задач без истории → используйте SQS, но это за рамками курса.

На практике всё чаще выбирают Kafka для потоков и RabbitMQ для запросов, работающих в связке.

contents

Заключение: будущее систем очередей. Границы между категориями стираются: Kafka внедряет очередь (KIP-932 — queue semantics), RabbitMQ улучшает персистентность, появляются системы вроде Pulsar и Redpanda. Однако фундаментальные различия (логирование vs маршрутизация) останутся. Для инженера важно не зацикливаться на одном инструменте, а понимать аксиомы распределённых систем: CAP-теорему и семантику доставки. Выбор должен основываться на требованиях бизнеса, а не на хайпе. Используйте сильные стороны каждого брокера и не бойтесь гибридных архитектур .