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







