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

Наша компанія займається розробкою, підтримкою та обслуговуванням сайтів будь-якої складності. Від простих односторінкових сайтів до масштабних кластерних систем, побудованих на мікро сервісах. Досвід розробників підтверджено сертифікатами від вендорів.

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

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

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

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

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

Етапи розробки

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

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

Реалізація 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.

Послуги бекенд-розробки: production-grade надійність

На production-сервері о 3:14 ночі черга Laravel Jobs перестала оброблятися — 40 000 необроблених завдань у Redis. Причина: worker упав через memory leak у статичній змінній Eloquent observer, supervisor не перезапустив через misconfigured stopwaitsecs. Ми розбирали такий інцидент на проекті з 500 RPS: діагностика 4 години, фікс — 20 хвилин. Щоб ви не втрачали гроші, пропонуємо послуги бекенд-розробки з акцентом на production-grade надійність — 10+ років досвіду, 50+ проектів, 5 років на ринку. Оцінимо ваш проект за 2 дні.

Які проблеми вирішуємо

N+1 запити: головний вбивця швидкості

N+1 — найпоширеніша причина повільних сторінок у Laravel-додатках. Стандартна історія: сторінка працювала нормально на dev з 10 записами, на production з 10 000 — 8-секундне завантаження.

Laravel Debugbar у dev-оточенні показує кількість запитів. Більше 20 — сигнал для audit.

Model::preventLazyLoading(! app()->isProduction());

Telescope для профілювання: логує всі запити, jobs, mail, notifications з деталізацією. Після впровадження eager loading час завантаження сторінки падає з 8 с до 0.3 с — у 27 разів.

Memory leak у статичних змінних

У Laravel Octane або Swoole додаток тримається в пам’яті між запитами. Статичні змінні не скидаються — призводять до неконтрольованого росту пам’яті. Використовуємо defer-функції та контейнерні біндинги для коректного скидання стану.

Неправильний connection pool

Rails, Laravel, Django відкривають нове з'єднання PostgreSQL на кожен PHP/Python процес. 100 воркерів — 100 з'єднань. PostgreSQL деградує від 200+ активних з'єднань через overhead на управління.

PgBouncer у transaction pooling: 1000 воркерів → 20–50 реальних з'єднань. Це знижує latency на 40% та зменшує витрати на хостинг на 30% — при середній вартості хостингу $2,000/міс економить $600/міс. GIN-індекс для JSONB до 100 разів швидший за B-tree при пошуку.

Як Octane справляється з високим навантаженням?

Laravel Octane (RoadRunner або Swoole) прибирає overhead bootstrap на кожен HTTP-запит. Приріст: 3–8x на синтетичних бенчмарках, 2–4x на реальних додатках. Важливо: не зберігати стан у статичних змінних — застосовуємо це на проектах >1000 RPS.

Як PostgreSQL допомагає уникнути повільних запитів?

Використовуємо composite indexes для WHERE + ORDER BY, partial indexes для фільтрів з високою селективністю, GIN-індекси для JSONB та full-text search. to_tsvector + GIN замість LIKE '%query%' — запобігає seq scan навіть на мільйонах записів. Аналізуємо плани через EXPLAIN ANALYZE та pg_stat_statements.

Як обрати стек для вашого проекту?

Стек Коли використовувати
Laravel + Octane CRUD, бізнес-логіка, REST/GraphQL API, адмінки
Node.js (Fastify) Realtime WebSocket, streaming, serverless, висока I/O concurrency
Go Високонавантажені мікросервіси (>10k RPS), gRPC, DevOps-інструменти
Django + DRF ML-пайплайни, інтеграція з AI, складна обробка даних
Ruby on Rails Швидкий MVP з багатим екосистемою гемів

Node.js виправданий для realtime: Laravel публікує події в Redis Pub/Sub, Node.js підписується та транслює клієнтам. Go — для goroutines (10k з'єднань на сервер — норма), але розробка повільніша, ніж Laravel.

Чому Redis критичний для продуктивності?

Redis виконує кілька ролей:

Роль Деталі
Кеш Кешування результатів важких запитів, фрагментів HTML
Черги Backend для Laravel Queue / Celery
Session store Distributed sessions в multi-instance оточенні
Pub/Sub Realtime події між сервісами
Rate limiting Sliding window counters для API throttling
Leaderboards Sorted Sets для рейтингів

Redis Cluster для горизонтального масштабування, Sentinel для автоматичного failover. Замовте консультацію щодо оптимізації Redis для вашого проекту.

Що входить в роботу під ключ

  • Архітектурне проектування (документація API, схема БД, діаграма сервісів)
  • Реалізація за узгодженим ТЗ з code review
  • Налаштування CI/CD (GitHub Actions, Docker), моніторингу (Sentry, Grafana), алертингу
  • Навантажувальне тестування (k6, wrk) зі звітом
  • Передача вихідних кодів, доступів, інструкція з деплою
  • Навчання команди замовника (2–3 сесії)
  • Гарантійна підтримка 1 місяць після здачі

Орієнтири по термінах

Задача Термін
REST API для мобільного/SPA (середня складність) 6–12 тижнів
Backend зі складною бізнес-логікою + інтеграції 12–20 тижнів
Високонавантажений сервіс на Go 8–16 тижнів
Міграція legacy PHP на Laravel 16–32 тижні

Вартість розраховується індивідуально після аналізу вимог до навантаження, інтеграцій та бізнес-логіки. Зв'яжіться з нами для безкоштовного аудиту вашого поточного backend — отримайте план оптимізації за 2 дні. Замовте консультацію та дізнайтеся, як знизити витрати на інфраструктуру на 30% без втрати продуктивності.