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, мониторинг и чеклист); общий выбор брокера — в отдельной статье.
В статье
- Зачем streams, если есть очереди
- Что такое stream: устройство
- Как объявлять и читать
- Offset и publish confirms
- Super Streams: партиционирование
- Streams против Kafka
- Можно ли заменить Kafka
- Где Streams кусаются
- Выводы
Зачем 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
Что такое 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", число, timestampOffset и 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
// объявляем 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 развиваются, детали меняются между релизами.
Комментарии