Реализация 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 рабочих дней в зависимости от сложности предметной области. Свяжитесь с нами — мы бесплатно оценим вашу архитектуру и предложим план внедрения.







