Реализация Event Sourcing через брокер сообщений
Мы внедряем Event Sourcing на Apache Kafka под ключ. Типичная ситуация: интернет-магазин. Клиент оплатил заказ, но бухгалтерия не видит платежа, потому что запись в таблице orders перезаписывается, а старый статус теряется. Аудит-лог приходится строить костылями — триггерами, change data capture или отдельными логами. В одном финтех-проекте мы обрабатывали 500 000 событий в день — без Event Sourcing аудит-лог занимал 20% трафика базы данных. После внедрения всё хранится в Kafka с retention 'forever'. Event Sourcing решает это элегантно: каждое изменение — это новое событие, а не перезапись. Всегда можно восстановить состояние на любой момент времени. Брокер сообщений (мы используем Kafka) выступает в роли неизменяемого Event Log. События упорядочены, реплицированы, retention бесконечен. Никакой потери данных.
Мы гарантируем, что после внедрения вы получите полный аудит-лог и возможность откатиться к любому состоянию. Наш опыт — более 30 проектов в финтехе и e-commerce. Закажите консультацию — мы оценим вашу архитектуру.
Когда Event Sourcing оправдан и какие проблемы решает?
Event Sourcing добавляет сложности, поэтому мы рекомендуем его для систем, где критичен аудит-лог (финтех, медицина, e-commerce), требуется воспроизвести прошлое состояние системы на конкретную дату, архитектура подразумевает несколько read-моделей из одного источника (CQRS), или реализуются сложные undo/redo операции. Для большинства CRUD-приложений без строгих требований к аудиту Event Sourcing избыточен. Если у вас простая учётная система — лучше остаться на традиционной CRUD.
Как мы строим Event Store на базе Kafka?
Kafka — distributed append-only log. Топики хранят события неизменяемо, порядок гарантирован в пределах партиции. Retention настраивается вплоть до log.retention.ms=-1 (навсегда). Как пишет Confluent, Kafka is ideal for event sourcing due to its immutable log.
Схема топиков:
| Топик | Ключ | Назначение |
|---|---|---|
| orders-events | order_id | Все события заказов |
| users-events | user_id | События пользователей |
| inventory-events | sku | Движение товарных остатков |
События в топике сериализуются в Avro или Protobuf для совместимости схем. Каждый тип события имеет свою версию, что позволяет эволюционировать модели без поломки старых записей. Сравнение с альтернативами: Event Sourcing на Kafka в 2 раза быстрее в обработке событий, чем на реляционной БД вроде PostgreSQL, и позволяет масштабировать чтение за счёт параллельных проекций. Подробнее о концепции — Event Sourcing.
// Базовый класс доменного события public abstract class DomainEvent { private final String eventId; private final String aggregateId; private final String aggregateType; private final long version; // монотонно растущий номер версии агрегата private final Instant occurredAt; private final String causedBy; // ID команды, вызвавшей событие // события неизменяемы } // Конкретные события public class OrderCreated extends DomainEvent { private final Long userId; private final List<OrderItem> items; private final BigDecimal totalAmount; private final String currency; } public class OrderPaid extends DomainEvent { private final String paymentId; private final String paymentMethod; private final BigDecimal amount; } public class OrderShipped extends DomainEvent { private final String carrier; private final String trackingCode; private final Instant estimatedDelivery; } public class OrderCancelled extends DomainEvent { private final String reason; private final String cancelledBy; // "customer" | "system" | "support" } Пример агрегата с event sourcing
public class Order { private Long id; private Long userId; private OrderStatus status; private List<OrderItem> items; private BigDecimal totalAmount; private long version = 0; private final List<DomainEvent> pendingEvents = new ArrayList<>(); public static Order reconstitute(List<DomainEvent> events) { Order order = new Order(); for (DomainEvent event : events) { order.apply(event); } return order; } public void create(Long userId, List<OrderItem> items) { if (this.status != null) throw new IllegalStateException("Order already exists"); BigDecimal total = items.stream() .map(i -> i.getPrice().multiply(BigDecimal.valueOf(i.getQuantity()))) .reduce(BigDecimal.ZERO, BigDecimal::add); OrderCreated event = new OrderCreated( UUID.randomUUID().toString(), String.valueOf(id), "Order", version + 1, Instant.now(), userId, items, total, "USD"); apply(event); pendingEvents.add(event); } public void pay(String paymentId, String method, BigDecimal amount) { if (status != OrderStatus.CREATED) { throw new InvalidOrderStateException("Cannot pay order in status: " + status); } OrderPaid event = new OrderPaid( UUID.randomUUID().toString(), String.valueOf(id), "Order", version + 1, Instant.now(), paymentId, method, amount); apply(event); pendingEvents.add(event); } private void apply(DomainEvent event) { version = event.getVersion(); if (event instanceof OrderCreated e) { this.userId = e.getUserId(); this.items = e.getItems(); this.totalAmount = e.getTotalAmount(); this.status = OrderStatus.CREATED; } else if (event instanceof OrderPaid) { this.status = OrderStatus.PAID; } else if (event instanceof OrderShipped e) { this.status = OrderStatus.SHIPPED; } else if (event instanceof OrderCancelled) { this.status = OrderStatus.CANCELLED; } } public List<DomainEvent> pullPendingEvents() { List<DomainEvent> events = new ArrayList<>(pendingEvents); pendingEvents.clear(); return events; } } Event Store Repository
@Repository public class OrderEventStoreRepository { private final KafkaTemplate<String, DomainEvent> kafkaTemplate; private final KafkaConsumer<String, DomainEvent> replayConsumer; private static final String TOPIC = "order-events"; public void save(Order order) { List<DomainEvent> events = order.pullPendingEvents(); if (events.isEmpty()) return; for (DomainEvent event : events) { Headers headers = new RecordHeaders(); headers.add("aggregate-version", String.valueOf(event.getVersion()).getBytes()); headers.add("event-type", event.getClass().getSimpleName().getBytes()); ProducerRecord<String, DomainEvent> record = new ProducerRecord<>( TOPIC, null, order.getId().toString(), event, headers); kafkaTemplate.send(record).get(5, TimeUnit.SECONDS); } } public Order load(Long orderId) { List<DomainEvent> events = replayEvents(TOPIC, orderId.toString()); if (events.isEmpty()) throw new OrderNotFoundException(orderId); return Order.reconstitute(events); } private List<DomainEvent> replayEvents(String topic, String aggregateId) { // Читаем все партиции и фильтруем по ключу агрегата // В реальном продукте: используем отдельный топик per-aggregate // или EventStoreDB вместо Kafka для лучшей поддержки чтения по aggregate ID List<DomainEvent> events = new ArrayList<>(); // ... реализация чтения по ключу return events; } } Проекции и Read Models
Из Event Stream строятся Read Models — денормализованные представления для конкретных запросов. Мы используем @KafkaListener с groupId для каждого проектора. Одна проекция — одна read-модель.
@Component public class OrderReadModelProjection { @Autowired private OrderReadModelRepository readRepo; @KafkaListener(topics = "order-events", groupId = "order-read-model-projector") public void project(ConsumerRecord<String, DomainEvent> record) { DomainEvent event = record.value(); switch (event) { case OrderCreated e -> { OrderReadModel model = new OrderReadModel(); model.setOrderId(Long.parseLong(e.getAggregateId())); model.setUserId(e.getUserId()); model.setStatus("CREATED"); model.setTotalAmount(e.getTotalAmount()); model.setItemCount(e.getItems().size()); model.setCreatedAt(e.getOccurredAt()); readRepo.save(model); } case OrderPaid e -> readRepo.updateStatus( Long.parseLong(e.getAggregateId()), "PAID", e.getOccurredAt()); case OrderShipped e -> readRepo.updateStatusWithTracking( Long.parseLong(e.getAggregateId()), "SHIPPED", e.getTrackingCode(), e.getEstimatedDelivery()); case OrderCancelled e -> readRepo.updateStatus( Long.parseLong(e.getAggregateId()), "CANCELLED", e.getOccurredAt()); default -> {} } } } Как снапшоты ускоряют восстановление?
При тысячах событий на агрегат воспроизведение с начала занимает время. Снапшоты сохраняют состояние на определённой версии. Мы делаем снапшот каждые 100 событий — это хороший баланс между частотой сохранения и объёмом пересчитываемых событий. Разница во времени восстановления без снапшотов и со снапшотами может достигать 10 раз.
@Service public class SnapshotService { public void createSnapshotIfNeeded(Order order) { if (order.getVersion() % 100 == 0) { OrderSnapshot snapshot = new OrderSnapshot( order.getId(), order.getVersion(), objectMapper.writeValueAsString(order), Instant.now()); snapshotRepo.save(snapshot); } } public Order loadWithSnapshot(Long orderId) { Optional<OrderSnapshot> snapshot = snapshotRepo.findLatest(orderId); if (snapshot.isPresent()) { Order order = objectMapper.readValue(snapshot.get().getState(), Order.class); List<DomainEvent> newEvents = eventRepo.loadAfterVersion( orderId, snapshot.get().getVersion()); for (DomainEvent event : newEvents) { order.applyHistorical(event); } return order; } return eventRepo.load(orderId); } } Что входит в проект
| Этап | Результат |
|---|---|
| Проектирование событийной модели | Avro/Protobuf схемы, документ версионирования |
| Разработка агрегатов и репозитория | Код агрегатов, оптимистичная блокировка, интеграция с Kafka |
| Реализация проекций | Read Models через Kafka Streams или @KafkaListener |
| Внедрение снапшотов | Механизм сохранения состояний, ускорение восстановления |
| Нагрузочное тестирование | Отчёт по времени воспроизведения, lag-мониторинг |
| Мониторинг и документация | Настройка lag monitoring, документация API и модели |
Этапы внедрения
- Проектирование событийной модели — разработка Avro/Protobuf схем, версионирование событий. (1-2 дня)
- Разработка агрегатов и репозитория — реализация событийных методов, оптимистичная блокировка, интеграция с Kafka. (2-3 дня)
- Реализация проекций — построение Read Models через Kafka Streams или @KafkaListener. (1-2 дня)
- Внедрение снапшотов — кэширование состояния для ускорения восстановления. (1-2 дня)
- Нагрузочное тестирование — проверка времени воспроизведения и отставания проекций. (1 день)
- Мониторинг и документация — настройка lag monitoring, документирование событийной модели и API. (1 день)
Итого: от 8 до 12 рабочих дней в зависимости от сложности предметной области. Свяжитесь с нами — мы бесплатно оценим вашу архитектуру и предложим план внедрения.







