Streams в RabbitMQ: Kafka-подобный лог и можно ли заменить Kafka

Что такое streams в RabbitMQ — append-only лог с offset-чтением и репликацией, как их объявлять и читать (native stream-протокол и AMQP), super streams для партиционирования, честное сравнение с Kafka и разбор, где RabbitMQ Streams реально заменяет Kafka, а где нет

RabbitMQ привычно воспринимают как брокер очередей: сообщение доставлено — сообщение удалено. Но с версии 3.9 в нём есть принципиально другой тип — stream: append-only лог, который читается по offset и не удаляет сообщение после доставки. Это тот же принцип, что у Kafka, и отсюда неизбежный вопрос: если у меня уже есть RabbitMQ, можно ли не тащить Kafka?

Разберём предметно: как streams устроены, как с ними работать (native stream-протокол и через AMQP), что такое super streams, чем они отличаются от Kafka — и где реально заменяют её, а где нет. Весь код прогнан на запускаемом стенде — digital-cookbook → rabbitmq/streams: docker-compose с RabbitMQ 4.3.2, Go stream-клиент v1.8.1, все примеры этой статьи с ожидаемым выводом. Streams активно развиваются, поведение меняется между версиями — поэтому пиньте версию.

Пятая статья серии про высокодоступный RabbitMQ (предыдущие — quorum-очереди и failover, durability/DLQ, Federation и Shovel, мониторинг и чеклист); общий выбор брокера — в отдельной статье.

RabbitMQ Streams: append-only лог с offset-чтением против партиций Kafka

В статье

Зачем streams, если есть очереди

Classic и quorum-очереди — это destructive consume: потребитель забрал сообщение, ack — сообщение ушло из очереди. Это отлично для задач-команд («обработай платёж»), но плохо, когда нужно:

  • перечитать историю с начала (replay) — новый сервис хочет обработать всё, что уже было;
  • раздать одно и то же многим независимым потребителям без копирования в N очередей (fan-out по логу);
  • держать долгую ретенцию и большой объём, читая с произвольной позиции;
  • высокий throughput на запись/чтение последовательного лога.

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

flowchart LR subgraph queue["Classic / Quorum queue"] direction LR Q[("очередь")] -->|"ack удаляет"| C1["Потребитель"] end subgraph stream["Stream — append-only лог"] direction LR L["лог: 0 1 2 3 4 …"] L -->|"offset 2"| R1["Потребитель 1"] L -->|"offset 4"| R2["Потребитель 2"] L -->|"replay c 0"| R3["Новый сервис"] end

flowchart LR
    subgraph queue["Classic / Quorum queue"]
        direction LR
        Q[("очередь")] -->|"ack удаляет"| C1["Потребитель"]
    end
    subgraph stream["Stream — append-only лог"]
        direction LR
        L["лог: 0 1 2 3 4 …"]
        L -->|"offset 2"| R1["Потребитель 1"]
        L -->|"offset 4"| R2["Потребитель 2"]
        L -->|"replay c 0"| R3["Новый сервис"]
    end
Очередь: ack удаляет сообщение. Stream: лог не меняется, потребители читают независимо и могут перечитать с начала

Что такое stream: устройство

Stream — это реплицируемый append-only лог:

  • append-only. Сообщения только дописываются в конец; чтение не меняет лог. Один и тот же offset читают сколько угодно потребителей.
  • offset-based. Позиция — это число (0, 1, 2, …). Потребитель читает с first, last, next, конкретного offset или с временной метки.
  • репликация. У stream есть leader и реплики на других нодах кластера (как quorum-очередь, но через отдельный лог-движок Osiris, а не Raft), поэтому он вписывается в тему высокой доступности серии.
  • ретенция по времени/размеру. Лог не бесконечен: задаётся максимальный возраст (MaxAge), суммарный размер (MaxLengthBytes) и размер сегмента (MaxSegmentSizeBytes). Старые сегменты усекаются.
  • два пути доступа. Native stream-протокол (отдельный порт 5552, выделенные клиенты, максимальный throughput) и обычный AMQP 0.9.1 (x-queue-type: stream) — стрим виден как очередь особого типа.

Важно: stream — не замена обычным очередям, а другой инструмент рядом с ними в том же брокере. Команды — в очереди, события-лог — в stream.

Как объявлять и читать

Native-путь (Go stream-клиент). Подключение идёт на порт 5552, стрим объявляется с ретенцией:

import (
    "github.com/rabbitmq/rabbitmq-stream-go-client/pkg/amqp"
    "github.com/rabbitmq/rabbitmq-stream-go-client/pkg/stream"
)

// полный запускаемый вариант — в digital-cookbook (go/stream)
env, err := stream.NewEnvironment(stream.NewEnvironmentOptions().
    SetHost("localhost").SetPort(5552).
    SetUser("guest").SetPassword("guest"))

// ретенция: не больше 5 ГБ лога (ещё есть SetMaxAge / SetMaxSegmentSizeBytes)
err = env.DeclareStream("events", stream.NewStreamOptions().
    SetMaxLengthBytes(stream.ByteCapacity{}.GB(5)))

// producer: Send асинхронный и батчится (см. про confirms ниже)
producer, err := env.NewProducer("events", nil)
_ = producer.Send(amqp.NewMessage([]byte("hello")))

// consumer: читаем с начала лога, callback на каждое сообщение
consumer, err := env.NewConsumer("events",
    func(_ stream.ConsumerContext, m *amqp.Message) {
        fmt.Printf("%s\n", m.GetData())
    },
    stream.NewConsumerOptions().
        SetConsumerName("worker").
        SetOffset(stream.OffsetSpecification{}.First()))

Тот же stream доступен и через AMQP — удобно, если не хочется тянуть stream-клиент. Но есть две обязательные детали: x-queue-type: stream при объявлении и — при чтении — basic.qos (prefetch) плюс x-stream-offset. Без prefetch стрим через AMQP работать не будет.

_, err = ch.QueueDeclare("events", true, false, false, false,
    amqp.Table{"x-queue-type": "stream"})

// для стрима ОБЯЗАТЕЛЬНЫ prefetch и стартовый offset
_ = ch.Qos(100, 0, false)
msgs, err := ch.Consume("events", "", false, false, false, false,
    amqp.Table{"x-stream-offset": "first"}) // или "last", "next", число, timestamp

Offset и publish confirms

Offset-спецификации (с какой позиции читать): First() — с начала лога, Last() — с последнего чанка, Next() — только новые, Offset(n) — с конкретной позиции, Timestamp(t) — с момента времени, LastConsumed() — с сохранённого offset. Именно Timestamp и произвольный Offset дают «time-travel», которого нет у обычных очередей.

Offset потребителя можно хранить на сервере — по имени потребителя, без внешнего хранилища:

// offset хранят ПЕРИОДИЧЕСКИ прямо из handler'а, а не на каждое сообщение
// (частый StoreOffset бьёт по производительности — предупреждение из README клиента):
handler := func(ctx stream.ConsumerContext, m *amqp.Message) {
    process(m)
    if processed%1000 == 0 {
        _ = ctx.Consumer.StoreOffset() // сохранить позицию под именем consumer ("worker")
    }
}
// при рестарте — прочитать сохранённое и продолжить с него:
off, err := env.QueryOffset("worker", "events")
// consumer с SetOffset(stream.OffsetSpecification{}.Offset(off + 1))

Про запись важно знать: Send асинхронный и батчится. Если закрыть producer сразу после цикла Send, часть сообщений может не успеть флашнуться и потеряется. Для гарантий подписывайтесь на publish confirmations (NotifyPublishConfirmation) и дожидайтесь подтверждения — это стрим-аналог publisher confirms из статьи про durability.

Кроме того, streams умеют дедупликацию на стороне producer: по имени producer и возрастающему publishing ID брокер отбрасывает повторы и помнит максимальный ID через рестарт — ближайший аналог идемпотентного продюсера Kafka. Но это не сквозной exactly-once (read-process-write): полного EOS у streams нет — впрочем, и у Kafka EOS ограничен границами её кластера.

Super Streams: партиционирование

Один stream — это один лог с полным порядком, но его пропускная способность и параллелизм чтения ограничены одной партицией. Super stream — это логическая композиция из N обычных стримов-партиций (orders-0, orders-1, …), поверх которой клиент делает роутинг по ключу. Прямой аналог партиционированного топика Kafka.

flowchart LR P["Super stream producer"] -->|"hash(key) mod 3"| R{"orders"} R --> S0["orders-0"] R --> S1["orders-1"] R --> S2["orders-2"] S0 --> C["Super stream consumer"] S1 --> C S2 --> C

flowchart LR
    P["Super stream producer"] -->|"hash(key) mod 3"| R{"orders"}
    R --> S0["orders-0"]
    R --> S1["orders-1"]
    R --> S2["orders-2"]
    S0 --> C["Super stream consumer"]
    S1 --> C
    S2 --> C
Super stream: producer роутит по хешу ключа в одну из партиций-стримов; consumer читает со всех
// объявляем super stream с 3 партициями (создаст orders-0, orders-1, orders-2)
err = env.DeclareSuperStream("orders", stream.NewPartitionsOptions(3))

// роутинг по application property "key" — сообщения одного ключа идут в одну партицию
routing := stream.NewHashRoutingStrategy(func(m message.StreamMessage) string {
    return m.GetApplicationProperties()["key"].(string)
})
producer, err := env.NewSuperStreamProducer("orders",
    stream.NewSuperStreamProducerOptions(routing))

msg := amqp.NewMessage([]byte("order created"))
msg.ApplicationProperties = map[string]any{"key": "user-42"}
_ = producer.Send(msg) // уедет в партицию hash("user-42") % 3

Порядок гарантируется в пределах партиции (как в Kafka), а не глобально: сообщения одного ключа всегда в одном стриме и читаются по порядку. Super stream consumer читает со всех партиций сразу. Для конкурентных потребителей с сохранением порядка есть single active consumer — в группе активен один, остальные на подхвате. (Super streams и single active consumer для стримов появились в RabbitMQ 3.11; на 4.3.2 доступны из коробки.)

Из практики. На живом брокере super stream с 3 партициями и хеш-роутингом по 10 ключам распределил 60 сообщений как 12/30/18 — хеш ключей ложится по партициям неравномерно, ровно как в Kafka. Это нормально; равномерность даёт не число сообщений, а разнообразие ключей.

Streams против Kafka

Модель одна — реплицируемый append-only лог с offset-чтением. Различия — в партиционировании, экосистеме и семантике.

Критерий RabbitMQ Streams Kafka
Модель append-only лог, реплицируемый append-only лог, реплицируемый
Партиционирование super stream (композиция стримов, роутинг на клиенте) нативные партиции топика
Порядок в пределах партиции в пределах партиции
Consumer groups single active consumer + super stream нативные consumer groups с ребалансом
Хранение offset на сервере по имени потребителя топик __consumer_offsets
Протокол RabbitMQ stream protocol (+ доступ по AMQP) Kafka protocol
Log compaction нет есть
Exactly-once нет сквозного (есть дедуп по publishing ID) транзакции / EOS (в границах Kafka)
Экосистема RabbitMQ-тулинг; нет Connect/Streams/ksqlDB Kafka Connect, Kafka Streams, ksqlDB
Соседство с очередями стримы и AMQP-очереди в одном брокере только Kafka

Можно ли заменить Kafka

Честный ответ — зависит от того, за чем вы шли в Kafka.

Streams реально заменяют Kafka, когда:

  • у вас уже есть RabbitMQ, и не хочется тащить второй кластер ради durable-лога;
  • нужны replay, fan-out по логу, длительная ретенция, чтение с произвольной позиции/времени;
  • нагрузка — высокая, но не экстремальная; партиционирование через super streams покрывает параллелизм;
  • в системе живут и команды (очереди), и события (лог) — удобно держать их в одном брокере.

Kafka не заменить, когда завязка именно на её экосистему и семантику:

  • Kafka Connect — готовые source/sink-коннекторы (CDC из БД, выгрузки в S3/ClickHouse и т.п.);
  • Kafka Streams / ksqlDB — потоковая обработка и стейт прямо на брокерных данных;
  • log compaction — «последнее значение по ключу» (changelog-топики, compacted-состояние);
  • exactly-once сквозь read-process-write через транзакции;
  • экстремальный масштаб с сотнями партиций и устоявшимся Kafka-тулингом/командной экспертизой.

Коротко: как durable event log внутри RabbitMQ-центричной системы — да, streams заменяют Kafka и убирают лишний кластер. Как платформу потоковой обработки с экосистемой — нет, это по-прежнему Kafka.

Где Streams кусаются

Подводные камни, которые лучше знать заранее (часть — из реального прогона):

  • Send теряет данные при небрежном закрытии. Он асинхронный и батчится; закрыли producer до флаша — часть сообщений не записана. Ждите publish confirmations.
  • AMQP-чтение без basic.qos не работает. Стриму через AMQP обязателен prefetch и x-stream-offset — забыли, и потребитель «молчит».
  • Offset — по имени потребителя, а не consumer group. Это не Kafka-группы с автоматическим ребалансом; семантику «поделить партиции между инстансами» дают super stream + single active consumer, и это надо выстроить руками.
  • Нет log compaction. «Последнее значение по ключу» на стримах не смоделировать — только полная ретенция по времени/размеру.
  • Ретенция обязательна. Без MaxAge/MaxLengthBytes лог растёт, пока не кончится диск. Задавайте пределы сразу.
  • Деплой: advertised_host. Клиент подключается к 5552, а дальше брокер сообщает ему хост партиции-лидера. Если stream.advertised_host/advertised_port настроены неверно (типичная засада за NAT/в контейнерах), соединение к стриму отвалится, хотя порт открыт.
  • Зрелость клиентов разнится по языкам. Native stream-клиенты живее всего для Go/Java/.NET; для остального часто остаётся путь через AMQP.

Выводы

  • Stream — это append-only лог внутри RabbitMQ: offset-чтение, неразрушающее потребление, реплики, ретенция. Отдельный инструмент рядом с очередями, а не замена им.
  • Работать можно и через native stream-протокол (максимальный throughput, super streams, серверные offset), и через AMQP (x-queue-type: stream + обязательный basic.qos).
  • Super streams дают партиционирование и порядок-в-пределах-партиции — прямой аналог партиций Kafka.
  • Заменить Kafka streams могут как durable event log в RabbitMQ-центричной инфраструктуре; не могут — там, где нужна экосистема Kafka (Connect, Streams/ksqlDB, log compaction, exactly-once, экстремальный масштаб).
  • Пиньте версию RabbitMQ и клиента: streams развиваются, детали меняются между релизами.

Документация и первоисточники

Обсуждение в Telegram

Присоединиться →

Комментарии