Перейти к содержанию

Spring & Kafka

Как это выглядит в Spring-мире

В Spring под "транзакциями" легко смешать три разных уровня:

  1. локальная транзакция внутри одного ресурса;
  2. настоящий distributed transaction через JTA/XA;
  3. координация без XA через события, Outbox и Saga.

Это важно разделять, потому что обычный:

@Transactional

решает в первую очередь локальную транзакцию внутри ресурса, которым управляет текущий transaction manager.

Он не делает автоматически:

  • общий commit между микросервисами;
  • атомарность между БД и Kafka;
  • глобальный rollback across services.

@Transactional: где его границы

Spring обычно использует transaction manager для:

  • открытия локальной транзакции;
  • commit / rollback;
  • привязки transactional context к текущему потоку исполнения.

Это хорошо работает для:

  • одной БД;
  • одного JPA persistence context;
  • одного JDBC resource manager.

Но если бизнес-метод делает так:

save order in DB
publish event to Kafka
call remote payment service

то одного @Transactional уже недостаточно, чтобы сделать все это одной честной транзакцией.


JtaTransactionManager и XA

Spring умеет работать с JTA через JtaTransactionManager.

Идея:

  • если есть XA-capable ресурсы;
  • и есть JTA transaction manager;
  • Spring может участвовать в distributed transaction сценарии.

Это ближе к enterprise / controlled environment use case, чем к typical microservices default.

Почему редко default:

  • XA support нужен у всех участников;
  • coordination тяжелее;
  • latency и failure modes сложнее;
  • внешние HTTP API и Kafka обычно не хочется тащить в XA-модель.

Поэтому в современных сервисах на Spring намного чаще встречается:

  • локальная транзакция;
  • Outbox;
  • Kafka;
  • Saga / asynchronous coordination.

@TransactionalEventListener

Spring поддерживает transaction-bound events через:

@TransactionalEventListener

Частая идея:

  • в локальной транзакции изменили данные;
  • после commit хотим дернуть какой-то listener.

Это полезно, когда нужно:

  • запустить локальную post-commit логику;
  • не выполнять side effect до успешного commit.

Но важно понимать границы:

  • listener не делает запись в Kafka атомарной вместе с БД;
  • если процесс упал после commit, но до надежной внешней публикации, событие можно потерять;
  • значит @TransactionalEventListener(AFTER_COMMIT) сам по себе не заменяет Outbox.

Его разумное использование:

  • после commit инициировать безопасную локальную реакцию;
  • или писать в Outbox / trigger relay logic;
  • но не считать это magic solution для dual write.

Как выглядит Outbox в Spring

Типичный паттерн в Spring-сервисе:

  1. внутри @Transactional сохранить бизнес-данные;
  2. там же вставить запись в outbox;
  3. закоммитить локальную транзакцию;
  4. отдельным relay-процессом отправить событие в Kafka.

Схема:

Spring service method
  @Transactional
    save aggregate
    insert outbox row
  COMMIT

relay
  -> read outbox
  -> publish to Kafka

Это как раз честный способ избежать:

DB commit succeeded
Kafka publish lost

Relay в Spring

Если делать polling publisher внутри Spring-приложения, обычно нужен отдельный компонент, который:

  • периодически читает пачку outbox rows;
  • публикует их в Kafka;
  • обновляет статус;
  • делает retry с backoff;
  • безопасно работает в нескольких инстансах.

Здесь сразу возникают практические вопросы:

  • как брать записи конкурентно;
  • как не отправить одно и то же слишком много раз;
  • как делать batching;
  • как обрабатывать publish success but status update failure.

Это еще одна причина, почему Outbox надо воспринимать как полноценный architectural pattern, а не как "еще одна таблица".


Debezium и Kafka

Если стек уже Kafka-centric, часто делают так:

  • приложение пишет только в свою БД и в outbox;
  • Debezium читает WAL/binlog;
  • Debezium публикует изменения outbox в Kafka.

Плюсы:

  • меньше кода в приложении;
  • логичное разделение ответственности;
  • хорошее сочетание с event-driven архитектурой.

Минусы:

  • нужна инфраструктура Kafka Connect / Debezium;
  • выше operational complexity;
  • не всегда оправдано для маленьких систем.

Saga на Spring + Kafka

Когда Saga строится поверх Kafka, Spring-сервис обычно делает такие роли:

  • принимает команду или событие;
  • выполняет локальную транзакцию;
  • пишет следующее событие через Outbox;
  • потребляет события других сервисов;
  • запускает следующий шаг или компенсацию.

Это может быть:

  • choreography-based flow через Kafka topics;
  • orchestration-based flow, где оркестратор тоже общается через Kafka и команды/события.

Критически важно:

  • стабильный event_id;
  • correlation_id;
  • saga_id;
  • идемпотентная обработка consumer'ов;
  • явные business statuses.

Exactly-once: где люди чаще всего ошибаются

Kafka может дать сильные гарантии внутри правильно построенного Kafka pipeline.

Но в Spring-бизнес-сервисе важно помнить:

  • запись в локальную БД и внешний side effect — уже не "бесплатный exactly-once";
  • даже если Kafka доставила аккуратно, consumer может упасть после локальной обработки;
  • Outbox снижает риск потери исходящего события, но не отменяет дубликаты;
  • именно поэтому Inbox, unique constraints и idempotent handlers остаются обязательными.

Практические правила для Spring-сервиса

  • Не пытайся делать межсервисную координацию обычным @Transactional.
  • Если нужно надежно публиковать событие после commit локальной БД, думай про Outbox.
  • Если процесс длинный и состоит из многих сервисов, думай про Saga.
  • @TransactionalEventListener — полезный инструмент, но не замена Outbox.
  • Закладывайся на повторную доставку и делай handlers идемпотентными.
  • Храни event_id, correlation_id, saga_id.

Связь с уже существующими заметками

По Spring-транзакциям и XA у тебя уже есть отдельная черновая заметка:

Этот файл лучше воспринимать как углубление в Spring-specific детали, а текущий раздел — как общую карту темы distributed transactions, Saga и Outbox.


Что важно запомнить

  • @Transactional отлично решает локальную транзакцию, но не распределенную координацию между сервисами.
  • Для XA/JTA в Spring есть отдельная модель через JtaTransactionManager, но это тяжелый и не всегда уместный путь.
  • В микросервисах на Spring чаще побеждает связка: local transaction + Outbox + Kafka + idempotent consumers + Saga where needed.
  • @TransactionalEventListener(AFTER_COMMIT) полезен, но сам по себе не решает dual write problem.