Event Sourcing на Kafka: впровадження під ключ

Реалізація Event Sourcing через брокер повідомлень

Розробка та обслуговування будь-яких видів сайтів:

Інформаційні сайти або веб-програми
Сайти візитки, landing page, корпоративні сайти, онлайн каталоги, квіз, промо-сайти, блоги, ресурси новин, інформаційні портали, форуми, агрегатори
Сайти або веб-програми електронної комерції
Інтернет-магазини, B2B-портали, маркетплейси, онлайн-обмінники, кешбек-сайти, біржі, дропшиппінг-платформи, парсери товарів
Веб-програми для управління бізнес-процесами
CRM-системи, ERP-системи, корпоративні портали, системи управління виробництвом, парсери інформації
Сайти або веб-програми електронних послуг
Дошки оголошень, онлайн-школи, онлайн-кінотеатри, конструктори сайтів, портали надання електронних послуг, відеохостинги, тематичні портали

Це лише деякі з технічних типів сайтів, з якими ми працюємо, і кожен із них може мати свої специфічні особливості та функціональність, а також бути адаптованим під конкретні потреби та цілі клієнта.

Послуги, які ми пропонуємо
Показано 1 з 1Усі 2062 послуг
Event Sourcing на Kafka: впровадження під ключ
Складний
~1-2 тижні

Наші компетенції:

Часті запитання

Останні роботи

  • image_website-b2b-advance_0.webp
    Розробка сайту компанії B2B ADVANCE
    1418
  • image_web-applications_feedme_466_0.webp
    Розробка веб-додатків для компанії FEEDME
    1286
  • image_websites_belfingroup_462_0.webp
    Розробка веб-сайту для компанії БЕЛФІНГРУП
    983
  • image_ecommerce_furnoro_435_0.webp
    Розробка інтернет магазину для компанії FURNORO
    1243
  • image_crm_enviok_479_0.webp
    Розробка веб-додатків для компанії Enviok
    983
  • image_bitrix-bitrix-24-1c_fixper_448_0.webp
    Розробка веб-сайту для компанії ФІКСПЕР
    998

Реалізація 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 та моделі

Етапи впровадження

  1. Проектування подієвої моделі — розробка Avro/Protobuf схем, версіонування подій. (1-2 дні)
  2. Розробка агрегатів та репозиторію — реалізація подієвих методів, оптимістичне блокування, інтеграція з Kafka. (2-3 дні)
  3. Реалізація проекцій — побудова Read Models через Kafka Streams або @KafkaListener. (1-2 дні)
  4. Впровадження снапшотів — кешування стану для пришвидшення відновлення. (1-2 дні)
  5. Навантажувальне тестування — перевірка часу відтворення та відставання проекцій. (1 день)
  6. Моніторинг та документація — налаштування lag monitoring, документування подієвої моделі та API. (1 день)

Разом: від 8 до 12 робочих днів залежно від складності предметної області. Зв'яжіться з нами — ми безкоштовно оцінимо вашу архітектуру та запропонуємо план впровадження.

Приклад архітектури Event Sourcing на Kafka Типова архітектура включає: івент-агрегатори (сервіси), які публікують події в Kafka; проектори, що будують read-моделі в базі даних; снапшот-сервіс, що періодично зберігає стан агрегатів; та моніторинг лагу через Kafka Lag Exporter.