В первой статье серии «транзакция» — это про изоляцию конкурентных чтений и атомарность изменений над данными: аномалии, уровни, BEGIN/COMMIT/ROLLBACK. В брокере сообщений это слово указывает на совершенно другой класс гарантий. Здесь не спрашивают «что увидит параллельный читатель посреди операции» — здесь спрашивают «дойдёт ли сообщение», «не потеряется ли оно при сбое» и «можно ли опубликовать пачку сообщений так, чтобы либо все они попали в лог, либо ни одно». RabbitMQ и Kafka отвечают на эти вопросы разными механизмами, и оба по пути показывают, где заканчивается их собственная транзакционность и начинается зона ответственности приложения.
В статье
- RabbitMQ: confirms вместо транзакций
- Kafka: транзакционный producer и exactly-once
- Границы EOS и outbox
- Карта: «транзакция» в брокере
- Что дальше
- Источники
RabbitMQ: confirms вместо транзакций
В AMQP 0-9-1, на котором построен RabbitMQ, транзакции существуют буквально — tx.select переводит канал в транзакционный режим, tx.commit подтверждает накопленные публикации и подтверждения (ack) сообщений, tx.rollback откатывает. Семантически это честный commit/rollback вокруг операций на канале. Проблема не в надёжности, а в цене: tx.commit — это синхронный round-trip к брокеру, и на время транзакции канал фактически заблокирован для следующей пачки операций. При публикации по одному сообщению за коммит throughput проседает на порядок по сравнению с обычной публикацией — плата за строгую транзакционную семантику, которую в очереди сообщений почти никогда не просят настолько буквально. В проде AMQP-транзакции RabbitMQ почти не используют именно поэтому.
Вместо них — publisher confirms: брокер асинхронно подтверждает каждую опубликованную запись — ack, если он принял на себя ответственность за сообщение, или nack, если что-то пошло не так. Публикации при этом не блокируют канал — можно отправить сразу пачку и параллельно вычитывать подтверждения из канала. Гарантия та же, что и от tx.commit, только стоит одного асинхронного канала подтверждений, а не синхронного round-trip на каждую публикацию.
Но что именно означает этот ack — место, где легко обмануться. Durable-очередь и persistent-сообщение — две разные настройки. durable=true при объявлении очереди говорит лишь о том, что рестарт переживёт сама очередь. Чтобы рестарт пережили и сообщения, публиковать их надо с DeliveryMode: amqp.Persistent. Без этого сообщение остаётся transient: ack придёт (брокер его действительно принял), очередь после рестарта на месте — а сообщений в ней нет.
// durable=true — переживёт рестарт САМА ОЧЕРЕДЬ
q, _ := ch.QueueDeclare("tx-demo", true, false, false, false, nil)
// включаем confirm mode на канале
if err := ch.Confirm(false); err != nil {
fmt.Println("confirm mode:", err)
return
}
confirms := ch.NotifyPublish(make(chan amqp.Confirmation, 200))
const n = 100
for i := 0; i < n; i++ {
// DeliveryMode: Persistent — чтобы рестарт пережили и СООБЩЕНИЯ.
// Без него ack придёт, но сообщение останется только в памяти.
if err := ch.PublishWithContext(ctx, "", q.Name, false, false, amqp.Publishing{
DeliveryMode: amqp.Persistent,
Body: []byte(fmt.Sprintf("msg-%d", i)),
}); err != nil {
fmt.Println("publish:", err)
return
}
}
acked, nacked := 0, 0
for i := 0; i < n; i++ {
c := <-confirms
if c.Ack {
acked++
} else {
nacked++
}
}Живой прогон на RabbitMQ 4: опубликовано 100, ack=100 nack=0 — каждая публикация подтверждена брокером, ни одного nack. Это и есть та надёжность, ради которой в старых руководствах предлагали AMQP-транзакции: клиент сразу узнаёт, принял брокер сообщение или нет, и может повторить публикацию.
А чтобы ack не читался шире, чем следует, стенд публикует две пачки по 100 сообщений в две одинаково durable-очереди: одну с DeliveryMode: amqp.Persistent, другую — без него. Обе получают ack=100. Разница видна только после рестарта брокера:
tx-demo (persistent): опубликовано 100, ack=100 nack=0
tx-demo-transient (transient): опубликовано 100, ack=100 nack=0
после рестарта: tx-demo осталось 100 из 100
после рестарта: tx-demo-transient осталось 0 из 100Обе очереди durable, обе публикации подтверждены — но пережили рестарт только persistent-сообщения. Так что ack означает «брокер принял на себя ответственность в рамках выбранного delivery mode», а не «данные на диске». Хотите durability — просите её явно, обеими настройками сразу: durable-очередь и persistent-публикация. И даже это не абсолют: persistent-сообщение попадает в буфер и сбрасывается на диск асинхронно, так что окно потери при жёстком отказе узла остаётся — закрывают его уже репликацией (quorum queues), а не одним флагом.
На стороне потребителя транзакции RabbitMQ вообще не участвуют — там своя модель: manual ack (basic.ack/basic.nack), at-least-once доставка и обязанность обработчика быть идемпотентным, потому что при переподключении или падении консьюмера до отправки ack брокер повторно доставит то же сообщение. Это отдельная гарантия, ортогональная publisher confirms, и её стоит держать в голове отдельно, не путая с «транзакционностью» публикации.
Отдельная заметка про стенд, если будете воспроизводить: в Docker на WSL RabbitMQ падает при старте на правах доступа к .erlang.cookie (eacces), пока контейнер не запущен сразу под пользователем rabbitmq — иначе Erlang успевает создать cookie от root, и брокер, работающий под uid 999, не может её прочитать. В compose стенда это учтено (user: rabbitmq), разбор — в README. К сути статьи отношения не имеет, но сэкономит полчаса, если поднимаете RabbitMQ в WSL самостоятельно.
Kafka: транзакционный producer и exactly-once
Kafka не пытается имитировать транзакции реляционной БД — вместо этого она даёт транзакционную атомарность публикации в лог: набор записей, отправленных в рамках одной транзакции, либо весь становится видимым читателям, либо весь остаётся невидимым, если транзакция прервана. Фундамент — idempotent producer: каждому продюсеру присваивается PID и epoch, каждой записи — порядковый номер в партиции, и брокер отбрасывает дубли при ретраях на уровне партиции. Транзакции строятся поверх этого механизма и добавляют координацию через transactional.id: он одновременно включает транзакционность и idempotent producer, а после падения продюсера с тем же transactional.id позволяет транзакционному координатору брокера завершить или откатить незакрытую предыдущую транзакцию раньше, чем начнётся новая — это то, что защищает от зомби-продюсеров.
Жизненный цикл транзакции на стороне producer: initTransactions() (один раз при старте — регистрирует transactional.id у координатора), beginTransaction(), серия send(), и в конце commitTransaction() либо abortTransaction(). Ключевая проверка — на стороне consumer: с isolation.level=read_committed записи прерванной транзакции не попадают в результат poll/PollFetches вообще, как будто их не было; с дефолтным read_uncommitted consumer увидел бы их все, включая заведомо отменённые.
prod, err := kgo.NewClient(
kgo.SeedBrokers(broker),
kgo.TransactionalID("tx-demo-producer"), // включает транзакции + idempotent producer
kgo.DefaultProduceTopic(topic),
kgo.AllowAutoTopicCreation(),
)
produceTx := func(prefix string, commit bool) {
if err := prod.BeginTransaction(); err != nil {
fmt.Println("begin tx:", err)
return
}
for i := 0; i < 10; i++ {
prod.Produce(ctx, &kgo.Record{Value: []byte(fmt.Sprintf("%s-%d", prefix, i))}, nil)
}
_ = prod.Flush(ctx)
kind := kgo.TryCommit
if !commit {
kind = kgo.TryAbort
}
if err := prod.EndTransaction(ctx, kind); err != nil {
fmt.Println("end tx:", err)
}
}
produceTx("aborted", false) // 10 записей в прерванной транзакции
produceTx("committed", true) // 10 записей в закоммиченной
cons, err := kgo.NewClient(
kgo.SeedBrokers(broker),
kgo.ConsumeTopics(topic),
kgo.FetchIsolationLevel(kgo.ReadCommitted()), // ключевое: не видим aborted
kgo.ConsumeResetOffset(kgo.NewOffset().AtStart()),
)Живой прогон на Apache Kafka 3.9.0: продюсер отправил 10 записей в прерванной транзакции и 10 в закоммиченной; read_committed consumer увидел committed=10, aborted=0 — прерванная транзакция полностью невидима, как будто её не отправляли.
Java (kafka-clients) — тот же жизненный цикл, но это и есть каноничный вид транзакционного API Kafka: именно на нём документация и KIP-98 формулируют семантику, а остальные клиенты (franz-go, confluent-kafka-go) её воспроизводят.
Properties pp = new Properties();
pp.put("bootstrap.servers", broker);
pp.put("transactional.id", "tx-demo-java");
pp.put("enable.idempotence", true);
try (KafkaProducer<String, String> producer = new KafkaProducer<>(pp)) {
producer.initTransactions();
produceTx(producer, topic, "aborted", false); // прерываем — записи не должны быть видны
produceTx(producer, topic, "committed", true); // коммитим — видны
}
Properties cp = new Properties();
cp.put("bootstrap.servers", broker);
cp.put("group.id", "tx-demo-java-consumer");
cp.put("auto.offset.reset", "earliest");
cp.put("isolation.level", "read_committed");
static void produceTx(KafkaProducer<String, String> producer, String topic, String prefix, boolean commit) {
producer.beginTransaction();
for (int i = 0; i < 10; i++) {
producer.send(new ProducerRecord<>(topic, prefix + "-" + i));
}
if (commit) {
producer.commitTransaction();
} else {
producer.abortTransaction();
}
}На идентичном сценарии Java даёт тот же результат: committed=10, aborted=0 — прерванная транзакция невидима read_committed consumer’у независимо от клиентской библиотеки.
Классический паттерн, который делает эту атомарность практически полезной, — consume-process-produce: сервис читает запись из входного топика, обрабатывает её и пишет результат в выходной топик, и всё это должно либо примениться целиком, либо не примениться вовсе — включая сдвиг оффсета входного топика. Если закоммитить выходную запись, но не сдвинуть оффсет (или наоборот), при падении между этими шагами сообщение обработается повторно или, наоборот, потеряется. Решение — sendOffsetsToTransaction: оффсет консьюмерской группы коммитится через продюсерскую транзакцию, в одной и той же транзакции с записью результата. Либо коммитятся и результат, и оффсет — цепочка продвинулась ровно на один шаг, либо не коммитится ничего — при повторном запуске сервис заново вычитает ту же входную запись и пересчитает результат. Это и есть exactly-once semantics в границах Kafka: не «эффект применится ровно один раз в мире», а «конвейер продвинется ровно на один шаг за одну успешную транзакцию, либо не продвинется вовсе».
Полный стенд — Go и Java — в digital-cookbook, transactions/brokers/ (проверено на RabbitMQ 4, Apache Kafka 3.9.0, PostgreSQL 18).
Границы EOS и outbox
Здесь стоит остановиться на честном ограничении: exactly-once semantics Kafka работает только внутри Kafka. Транзакция гарантирует атомарность между записями в разные топики и партиции и коммитом оффсетов той же consumer-группы — и ничего сверх этого. Она не покрывает запись в PostgreSQL, вызов внешнего HTTP API, отправку письма — любой сайд-эффект за пределами кластера Kafka происходит вне транзакционной границы. Распределённой транзакции, которая атомарно охватывала бы и PostgreSQL, и Kafka, не существует — 2PC между ними на практике не применяют: ни у одной из систем нет для этого готового промышленного координатора, и цена по латентности была бы неприемлемой даже если бы был.
Практическое решение — outbox pattern, который смыкает эту статью с транзакциями PostgreSQL из второй статьи серии. Идея простая: бизнес-изменение и запись о событии, которое из него следует, пишутся в одной транзакции той же СУБД — атомарность здесь обеспечивает не Kafka, а обычный ACID-движок реляционной БД. Отдельный процесс-релей затем читает неопубликованные строки из таблицы outbox, публикует каждую в Kafka и, только получив подтверждение публикации, помечает строку как published. Если публикация в Kafka упадёт — событие всё равно durable в БД (транзакция с бизнес-данными уже закоммичена независимо от брокера), и релей просто повторит попытку на следующем проходе. Событие никогда не теряется и никогда не публикуется без коммита бизнес-данных, потому что порядок жёсткий: сначала атомарная запись в БД, потом (возможно, с задержкой и повторами) публикация.
// бизнес-операция + событие — в ОДНОЙ транзакции: либо оба, либо ни одного
for i := 0; i < n; i++ {
tx, err := pool.Begin(ctx)
if err != nil {
continue
}
_, _ = tx.Exec(ctx, `INSERT INTO orders (item) VALUES ($1)`, fmt.Sprintf("item-%d", i))
_, _ = tx.Exec(ctx, `INSERT INTO outbox (payload) VALUES ($1)`, fmt.Sprintf("order-created-%d", i))
_ = tx.Commit(ctx)
}
// релей: читаем неопубликованные, публикуем в Kafka, помечаем published
rows, _ := pool.Query(ctx, `SELECT id, payload FROM outbox WHERE NOT published ORDER BY id`)
for rows.Next() {
var o ob
_ = rows.Scan(&o.id, &o.payload)
pending = append(pending, o)
}
rows.Close()
published := 0
for _, o := range pending {
if err := prod.ProduceSync(ctx, &kgo.Record{Value: []byte(o.payload)}).FirstErr(); err == nil {
mustExec(ctx, pool, `UPDATE outbox SET published = true WHERE id = $1`, o.id)
published++
}
}Живой прогон: 20 заказов + 20 событий записаны атомарно в одной БД-транзакции (PostgreSQL 18) — либо оба INSERT, либо ни одного, для каждого из 20 заказов; релей опубликовал 20 событий в Kafka; неопубликованных после прогона осталось 0. Outbox — это мост: он не делает Kafka частью транзакции БД и не делает БД частью транзакции Kafka, а просто гарантирует, что публикация в брокер физически не может произойти без предварительного, независимо durable, коммита бизнес-данных.
Карта: «транзакция» в брокере
| Система | Что значит «транзакция» / атомарность | Что НЕ покрывает |
|---|---|---|
| RabbitMQ | AMQP-транзакции существуют, но дороги; на практике — publisher confirms: асинхронный ack на каждую публикацию — «брокер принял ответственность в рамках delivery mode». Durability надо просить явно и дважды: durable-очередь и DeliveryMode: Persistent (иначе рестарт переживёт очередь, но не сообщения) |
Не изоляция чтений — на стороне consumer своя модель (manual ack, at-least-once, идемпотентность обработчика) |
| Kafka | Транзакционный producer + EOS: атомарная публикация в несколько топиков/партиций и атомарный коммит оффсетов consumer-группы через sendOffsetsToTransaction |
Не покрывает внешние сайд-эффекты — запись в БД, вызов API; граница гарантии — периметр кластера Kafka |
| outbox | Мост БД↔брокер: атомарность бизнес-данных и события обеспечивает ACID-транзакция БД, а не брокер; релей публикует независимо, с повтором при сбое | Не даёт мгновенной публикации — есть задержка между коммитом в БД и публикацией; не отменяет необходимость идемпотентности на стороне потребителя события |
Вывод для всей статьи: в брокерах транзакционность — это про доставку и про атомарность публикации, а не про то, что видит конкурентный читатель посреди чужой операции (это вопрос из мира БД, статьи #1–#3). Как только на сцену выходит внешняя система — БД, другой брокер, HTTP-сервис — распределённой транзакции, которая накрыла бы всё разом, ждать не стоит; вместо неё нужен мост вроде outbox, построенный на атомарности одной из систем и повторных попытках со стороны другой.
Что дальше
Финал серии — мульти-хранилищный стенд и карта выбора: один сценарий прогоняется по всем хранилищам сразу, и вся серия сводится в единую карту компромиссов.
Комментарии