Kafka: очередь сообщений и потоковая платформа — Часть 6

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

О 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, а также где на этом пути могут появиться дубликаты или потеряться данные. Потому что в распределённых системах «почти доставили» — это не статус, а начало расследования.