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.