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.

Услуги бэкенд-разработки: Laravel, Node.js, Go, Django, PostgreSQL

На production-сервере в 3:14 ночи очередь Laravel Jobs перестала обрабатываться. 40 000 необработанных задач в Redis. Причина: worker упал из-за memory leak в одном из Jobs (утечка через статическую переменную в Eloquent observer), supervisor не перезапустил его из-за misconfigured stopwaitsecs. Это не гипотетический сценарий — это вторник. Мы разбирали такой инцидент на проекте с нагрузкой 500 RPS: диагностика заняла 4 часа, фикс — 20 минут. Чтобы вы не теряли деньги на простоях, предлагаем услуги бэкенд-разработки с акцентом на production-grade надёжность. Оценим ваш проект за 2 дня.

Backend — это то, что работает когда никто не смотрит. Или не работает. Гарантируем, что у вас будет первый вариант.

Что мы делаем с первого дня правильно

Service Layer поверх Fat Controllers. Controller получает HTTP-запрос, валидирует его через Form Request, передаёт данные в Service, возвращает ответ. Бизнес-логика в Service, не в Controller. Это звучит банально, но большинство legacy-проектов — это контроллеры по 500 строк с SQL-запросами внутри.

Repository Pattern используем осторожно. Если вы просто оборачиваете Model::where(...) в метод репозитория — это бойлерплейт без пользы. Repository оправдан когда: нужно абстрагироваться от источника данных (БД + кеш + внешний API) или когда логика запросов достаточно сложна для изоляции.

Jobs, Events, Listeners. Всё, что можно сделать асинхронно — делаем асинхронно. Отправка email, генерация PDF, синхронизация с внешним API, пересчёт агрегатов — в Queue. Laravel Horizon для мониторинга очередей в Redis: видно throughput, failed jobs, время обработки по очередям.

Как Octane справляется с высокой нагрузкой

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

Что делать с N+1 запросами

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

Laravel Debugbar в dev-окружении показывает количество запросов на страницу. Более 20 запросов на одну страницу — сигнал для audit.

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

Telescope для профилирования в staging: логирует все запросы, jobs, mail, notifications с детализацией по времени. Цифры: после внедрения eager loading время загрузки страницы падает с 8 с до 0.3 с — в 27 раз.

PostgreSQL: индексы, которые реально нужны

PostgreSQL 14+ — основная БД на всех проектах. Используем связку PgBouncer + PostgreSQL. Опыт 10+ лет, более 50 backend-проектов, 5 лет на рынке.

Как PostgreSQL помогает избежать медленных запросов

Composite indexes для частых WHERE + ORDER BY. Если у вас WHERE user_id = ? AND status = ? ORDER BY created_at DESC — нужен (user_id, status, created_at DESC). Индекс по (user_id) отдельно плохо помогает с сортировкой.

Partial indexes. Если 95% запросов идут по WHERE status = 'active':

CREATE INDEX idx_orders_active ON orders (created_at DESC)
WHERE status = 'active';

Индекс маленький, быстрый, покрывает основную нагрузку.

GIN-индексы для JSONB и массивов. @> оператор без GIN-индекса — seq scan. С индексом — быстро даже на миллионах записей.

GIN для full-text search. to_tsvector + GIN вместо LIKE '%query%'. LIKE без индекса — всегда seq scan. С pg_trgm extension и gin_trgm_ops — поддержка LIKE с индексом, полезно для CRM-поиска по частичному совпадению.

Connection pooling: почему важнее чем кажется

Rails, Laravel, Django открывают новое соединение с PostgreSQL на каждый PHP/Python процесс. На 100 воркерах — 100 соединений. PostgreSQL начинает деградировать от 200–300 активных соединений — overhead на управление соединениями становится значительным.

PgBouncer — connection pooler перед PostgreSQL. Режим transaction pooling: соединение с PostgreSQL занято только на время транзакции, между запросами возвращается в пул. 1000 приложений-воркеров → 20–50 реальных соединений к PostgreSQL. Это снижает latency на 40% и уменьшает затраты на хостинг на 30%.

Node.js с Fastify: когда это лучше Laravel

Node.js оправдан для:

  • Realtime: WebSocket-серверы, Server-Sent Events, чат, live-обновления
  • Streaming: большие файлы, видео, данные потоком
  • High I/O concurrency: много параллельных запросов к внешним API без тяжёлой бизнес-логики
  • Serverless: Lambda/Cloud Functions — Node.js стартует быстрее PHP

Fastify вместо Express: в 2–3 раза быстрее на benchmarks, встроенная JSON Schema валидация, лучшая TypeScript поддержка, plugin-архитектура.

Типичная архитектура realtime: Laravel — основная бизнес-логика и REST API. Node.js + Socket.io или ws — WebSocket сервер. Laravel публикует события в Redis Pub/Sub, Node.js подписывается и транслирует клиентам. Это разделение позволяет масштабировать WebSocket-сервер независимо от основного приложения.

Go: микросервисы и высокая нагрузка

Go используем для:

  • Высоконагруженных микросервисов (> 10 000 RPS)
  • Фоновых воркеров с жёсткими требованиями к latency
  • Инструментов DevOps и CLI
  • gRPC-сервисов в микросервисной архитектуре

Goroutines — дешевле OS-потоков в тысячи раз. 10 000 конкурентных соединений на Go — норма на одном сервере.

Но Go — не волшебная таблетка. Разработка медленнее чем на Laravel: больше бойлерплейта, нет ORM уровня Eloquent, обработка ошибок через if err != nil везде. Оправдан только когда производительность — реальное требование, не предположение.

Django и Python backend

Django с DRF (Django REST Framework) — для задач где нужен Python: ML-пайплайны, обработка данных, интеграции с AI-инструментами.

Celery для фоновых задач — аналог Laravel Queue, но сложнее в конфигурации. Celery Beat для cron-задач.

Django ORM vs raw SQL: ORM удобен для CRUD. Для аналитических запросов с несколькими JOIN, оконными функциями и CTE — connection.execute() с raw SQL читаемее и предсказуемее.

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 на standalone установках.

Деплой и инфраструктура

Docker + docker-compose — стандарт для локальной разработки и production. Каждый сервис в контейнере: PHP-FPM/Octane, Nginx, PostgreSQL, Redis, Queue Worker, Scheduler.

CI/CD через GitHub Actions:

  1. Прогон тестов (PHPUnit / Pest, Vitest, Playwright)
  2. Сборка Docker-образа
  3. Push в Container Registry
  4. Deploy: docker pull → docker-compose up -d на сервере, или Kubernetes rolling update

Zero-downtime deploy для Laravel: php artisan down --secret=TOKEN не нужен при правильной настройке. Стратегия: новый контейнер стартует рядом со старым, Nginx переключает трафик после health check, старый контейнер останавливается.

Мониторинг: Sentry для exception tracking с alerting в Slack/Telegram. Grafana + Prometheus (или Grafana Cloud) для метрик: CPU, memory, request rate, queue depth, database connection count. Алерт на: error rate > 1%, p99 latency > 2s, queue depth > 1000 jobs.

Что входит в работу под ключ

  • Архитектурное проектирование (документация API, схема БД, диаграмма сервисов)
  • Реализация по согласованному ТЗ с code review
  • Настройка CI/CD, мониторинга, алертинга
  • Нагрузочное тестирование (k6, wrk) с отчётом
  • Передача исходников, доступов, инструкция по деплою
  • Обучение команды заказчика (2-3 сессии)
  • Гарантийная поддержка 1 месяц после сдачи

Ориентиры по срокам

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

Стоимость рассчитывается индивидуально после анализа требований к нагрузке, интеграциям и бизнес-логике. Типичный бюджет backend-проекта — от 500 000 до 2 000 000 рублей в зависимости от сложности. Свяжитесь с нами для бесплатного аудита вашего текущего backend — получите план оптимизации за 2 дня. Закажите консультацию.