Kafka: очередь сообщений и потоковая платформа — Часть 6
О Kafka часто говорят как о брокере сообщений.
Это не совсем ошибка, но и не полная картина. Примерно как назвать Kubernetes «штукой, которая запускает контейнеры»: технически верно, но где-то рядом уже грустит SRE.
Kafka действительно может работать как очередь. Однако её архитектура позволяет использовать один и тот же Topic сразу в двух режимах:
- как механизм распределения фоновых задач;
- как источник непрерывного потока событий.
Именно поэтому Kafka называют не просто message broker, а Event Streaming Platform.
Классическая очередь сообщений
Сначала вспомним, как обычно работает очередь.
Producer
|
v
Queue
|
v
Consumer
Producer создаёт задачу.
Consumer получает её и обрабатывает.
Например:
Сгенерировать отчёт
Обработать изображение
Отправить письмо
Пересчитать статистику
Главная идея:
Одну задачу должен выполнить один обработчик.
Если у нас несколько Consumer, очередь распределяет задачи между ними.
Task 1 -> Consumer A
Task 2 -> Consumer B
Task 3 -> Consumer C
Такой подход хорошо подходит для фоновых операций: сделать отчёт, пережевать картинку, отправить письмо и не заставлять пользователя смотреть на спиннер как на смысл жизни.
Kafka в роли очереди
Kafka может реализовать похожую модель с помощью Consumer Group.
Представим Topic:
video-processing
В него поступают задачи на обработку видео.
Video 1
Video 2
Video 3
Video 4
Есть группа обработчиков:
Consumer Group: transcoding-workers
Kafka распределяет Partition между ними.
Partition 0 -> Worker A
Partition 1 -> Worker B
Partition 2 -> Worker C
Каждое сообщение обрабатывается одним Consumer внутри группы.
Для внешнего наблюдателя это выглядит как обычная очередь задач.
Когда Kafka использовать как очередь
Такой режим подходит для асинхронной тяжёлой работы.
Например:
- конвертация видео;
- генерация документов;
- обработка фотографий;
- отправка уведомлений;
- расчёт рекомендаций;
- построение отчётов.
Пользователь не должен ждать завершения операции в HTTP-запросе.
Сервис просто создаёт событие:
{
"videoId": 15321,
"action": "transcode"
}
А Worker обработает его позже.
Kafka как буфер нагрузки
Представим, что Producer создаёт 100 000 сообщений в секунду.
Consumer успевает обрабатывать только 20 000.
Без промежуточного буфера сервис быстро окажется перегружен.
Kafka принимает входящий поток и сохраняет его в журнале.
Producer: 100 000 msg/s
|
v
Kafka
|
v
Consumer: 20 000 msg/s
Сообщения не теряются.
Они накапливаются внутри Topic, пока Consumer не сможет их обработать.
В этом сценарии Kafka работает как амортизатор между быстрым Producer и медленным Consumer.
Такой механизм называется backpressure buffering. Исходный материал описывает Kafka именно как средство развязки систем и поглощения всплесков нагрузки.
Пример с загрузкой видео
Представим видеохостинг.
Пользователь загрузил ролик.
Система должна:
- сохранить оригинал;
- сделать версию 1080p;
- сделать версию 720p;
- создать превью;
- проверить содержимое;
- извлечь аудиодорожку.
Выполнять всё синхронно нельзя.
Пользователь будет ждать несколько минут.
Вместо этого сервис создаёт события.
VideoUploaded
GeneratePreview
Transcode1080
Transcode720
ExtractAudio
Worker читают их из Kafka и выполняют задачи независимо.
При росте нагрузки достаточно добавить новых Consumer.
Ограничение очередной модели
Есть важный нюанс.
Kafka не является классической очередью задач во всех смыслах.
В обычной очереди сообщение часто удаляется сразу после обработки.
В Kafka оно остаётся до окончания срока хранения.
Offset 0 processed
Offset 1 processed
Offset 2 processed
Offset 3 waiting
Consumer лишь фиксирует свой Offset.
Поэтому Kafka лучше воспринимать как журнал, поверх которого можно построить поведение очереди.
Теперь рассмотрим потоковую модель
В режиме очереди сообщение воспринимается как задача:
Сделай что-нибудь один раз.
В потоковой модели сообщение воспринимается как факт:
Что-то произошло.
Например:
UserRegistered
OrderPaid
ProductViewed
FlightPriceChanged
Такое событие не обязательно должно быть «выполнено».
Разные системы могут наблюдать его и реагировать независимо.
Один поток — много подписчиков
Представим событие:
OrderPaid
Его одновременно хотят получить:
- сервис доставки;
- сервис аналитики;
- CRM;
- система уведомлений;
- антифрод;
- программа лояльности.
Каждый сервис создаёт свою Consumer Group.
OrderPaid
|
+---------------+---------------+
| | |
v v v
Delivery Analytics Notification
Каждая группа получает полную копию потока.
Ни один Consumer не мешает остальным.
Это и есть модель Publish/Subscribe.
Kafka как поток событий
В потоковой модели Topic становится непрерывной историей изменений.
12:00 UserRegistered
12:01 ProductViewed
12:02 ProductAddedToCart
12:03 OrderCreated
12:04 PaymentCompleted
Consumer читает этот поток и строит собственное представление данных.
Например:
- аналитика считает конверсию;
- рекомендации обновляют профиль пользователя;
- антифрод ищет подозрительные действия;
- CRM обновляет карточку клиента.
Каждый сервис интерпретирует одни и те же события по-своему.
Пример с рекламными кликами
Допустим, необходимо считать клики по рекламе в реальном времени.
Каждое нажатие создаёт событие:
{
"campaignId": 512,
"userId": 9812,
"timestamp": "2026-08-06T09:20:00Z"
}
Эти события постоянно поступают в Kafka.
Click
Click
Click
Click
Click
Consumer может:
- считать CTR;
- выявлять мошеннические клики;
- обновлять бюджеты кампаний;
- строить realtime-дашборд.
Здесь нет отдельной «задачи», которую нужно удалить после выполнения.
Есть непрерывный поток фактов.
Именно такой сценарий исходная статья называет «real-time river» — непрерывной рекой событий.
Главное отличие очереди от потока
Можно свести различие к одному вопросу.
Очередь
Кто должен выполнить эту задачу?
Обычно ответ:
Один Worker
Поток
Кто заинтересован в этом событии?
Ответ:
Любое количество независимых систем
Сравнение подходов
| Характеристика | Очередь | Поток |
|---|---|---|
| Основная сущность | Задача | Событие |
| Количество обработчиков | Обычно один | Один или много |
| Сообщение после чтения | Считается выполненным | Остаётся в журнале |
| Повторное чтение | Обычно ограничено | Поддерживается через Offset |
| Типичный сценарий | Фоновая обработка | Аналитика и реакция на события |
| Масштабирование | Несколько Worker | Несколько Consumer Group |
Kafka позволяет совместить оба подхода
Это одна из самых сильных сторон Kafka.
Один сервис может читать Topic как очередь задач.
Другой — как поток событий.
Представим Topic:
orders
Группа delivery-workers воспринимает каждое сообщение как задачу:
Создать доставку
Группа analytics воспринимает те же сообщения как события:
Обновить статистику заказов
Группа fraud-detection анализирует их как поток:
Проверить подозрительную активность
Все работают независимо.
orders
|
+---------------+---------------+
| | |
v v v
Delivery Analytics Fraud
Queue Stream Stream
Один Topic одновременно обслуживает разные модели использования.
Когда выбирать режим очереди
Используйте Kafka как очередь, когда:
- задачу должен обработать один
Worker; - операция выполняется асинхронно;
- нагрузку нужно распределить;
Producerработает быстрееConsumer;- важна возможность пережить всплеск трафика.
Типичные примеры:
Генерация PDF
Обработка видео
Отправка email
Фоновый импорт данных
Проверка файлов
Когда использовать поток
Потоковая модель подходит, если:
- одно событие нужно нескольким сервисам;
- важна история событий;
- данные обрабатываются в реальном времени;
- требуется агрегировать непрерывный поток;
Consumerдолжны читать независимо друг от друга.
Типичные примеры:
Клики пользователей
Изменения цен
Банковские транзакции
Телеметрия
Логи
События заказов
Kafka и Event-Driven Architecture
Kafka особенно хорошо раскрывается в событийной архитектуре.
Вместо команды:
Отправь пользователю письмо
система публикует факт:
UserRegistered
А сервис уведомлений сам решает, что делать.
UserRegistered
|
v
Email Service
|
v
Welcome Email
Появится новый сервис — например, программа лояльности — ему достаточно подписаться на то же событие.
Producer менять не нужно.
Команды и события — не одно и то же
Важно различать два типа сообщений.
Команда
Команда содержит намерение выполнить действие.
SendWelcomeEmail
Обычно у неё один конкретный получатель.
Событие
Событие описывает уже произошедший факт.
UserRegistered
Получателей может быть сколько угодно.
Для Kafka чаще предпочтительнее события, потому что они уменьшают связанность между системами.
Практическое правило
При проектировании Topic задайте себе вопрос:
Сообщение описывает действие, которое нужно выполнить, или факт, который уже произошёл?
Если действие:
GenerateInvoice
Kafka используется ближе к очередной модели.
Если факт:
InvoiceGenerated
Kafka используется как поток событий.
Итог
Kafka нельзя считать только очередью сообщений.
Она объединяет сразу несколько возможностей:
- распределённый журнал;
- асинхронную очередь;
- Publish/Subscribe;
- потоковую платформу;
- долговременное хранилище событий.
Её сила заключается не в том, что она идеально заменяет классические очереди.
Сила Kafka в том, что один и тот же поток данных можно читать разными способами и использовать для совершенно разных задач. Один сервис видит задачу, другой — аналитику, третий — повод поднять алерт.
Что дальше?
Мы уже разобрали:
TopicиPartition;ProducerиConsumer;Consumer Group;Offset;- append-only log;
- репликацию;
- причины высокой производительности;
- очередь и потоковую модель.
В следующей части соберём всё вместе и пройдём полный путь одного сообщения:
Producer
|
v
Partitioner
|
v
Leader Broker
|
v
Follower Replicas
|
v
Consumer Group
|
v
Offset Commit
Разберём, что происходит с момента вызова send() до подтверждения обработки Consumer, а также где на этом пути могут появиться дубликаты или потеряться данные. Потому что в распределённых системах «почти доставили» — это не статус, а начало расследования.