Привет, Хабр! Я Артём Борисов, Java-разработчик, в основном занимаюсь развитием микросервисов в команде РСХБ «Свои инвестиции». Представьте ситуацию: вы работаете с инвестиционными сделками, обрабатываете миллионы сделок в день, но все они обрабатываются один раз только ночью. А бизнес требует реального времени. Это была наша рутина, пока мы не внедрили Kafka Streams. В этой статье я расскажу о том, как мы трансформировали систему обработки сделок на фондовом рынке (SOFR) с batch-обработки на полноценную real-time систему, способную обрабатывать миллионы сделок в сутки.

Только в реальном времени

Наша система была классическим монолитом, где все сделки и заявки с различных финансовых рынков обрабатывались один раз в день — в 00:00 ночи. Процесс выглядел так: 

Торговый день → Накопление сделок → 00:00 → Выгрузка → Сопоставление отчетов

Это было неудобно: в течение дня нельзя было быстро реагировать на проблемы. Мы получали очереди на обработку в ночные часы. И я уже не говорю о задержках в выявлении ошибок, а это очень критичный момент в финансовой сфере. Бизнес потребовал изменений: они хотели, чтобы система обрабатывала каждую сделку в момент её поступления, но никак не следующей ночью.

Что мы сделали? Использовали Kafka Streams и GlobalKTable. Выбрали Kafka Streams по нескольким причинам — пропускная способность, встроенная обработка состояния (а значит не нужна отдельная база данных), масштабируемость, простота развёртывания (просто Java-приложение).

Что было в архитектуре решения? Ключевой компонент нашего решения — GlobalKTable:

Создаем универсальный метод для создания наших глобальных таблиц

private <K, V> GlobalKTable<K, V> createGlobalTable(

            String topicName,

            StreamsBuilder builder,

            Serde<K> keySerde,

            Serde<V> valueSerde

    ) {

             return builder.globalTable(

                topicName,

  Consumed.with(keySerde, valueSerde),

                Materialized.<K, V, KeyValueStore<Bytes, byte[]>>as("reference-data")

                        .withKeySerde(keySerde)

                        .withValueSerde(valueSerde)

        );

    }

Затем создаем метод под конкретный справочник

@Bean 
public GlobalKTable<String, TestData> testDataTable(StreamsBuilder builder) {

        return createGlobalTable(

                topics.getTestData(),

                builder,

                Serdes.String(),

                new JacksonJsonSerde<>(TestData.class)

        );

    }

Обогащаем сделку данными из справочника

KStream<String, Deal> deals = builder.stream("deals-topic");

KStream<String, EnrichedDeal> enrichedDeals = deals.join(

    referenceTable,

    (dealKey, deal) -> deal.getReferenceKey(),

    (deal, refData) -> enrichDeal(deal, refData)

);

enrichedDeals.to("enriched-deals-topic");

Почему GlobalKTable?

При знакомстве с Kafka Streams возникает закономерный вопрос: зачем использовать GlobalKTable, если существует обычный KTable? .

GlobalKTable работает иначе: каждый экземпляр приложения получает полную копию справочника и хранит её локально. Благодаря этому любое обогащение выполняется через локальный справочник без обращения к другим узлам.

Упрощённо различие выглядит следующим образом:

KTable

Deal -> repartition -> join -> результат

GlobalKTable

Deal -> локальный lookup -> join -> результат

В нашем случае справочники относительно небольшие, а скорость обработки сделок критически важна. Поэтому мы сознательно выбрали GlobalKTable: он загружает весь справочник локально на каждый экземпляр приложения. Обработка сделки перестала зависеть от производительности базы данных или доступности сторонних сервисов. Время доступа стало измеряться микросекундами.

Конечно, в GlobalKTable есть и ограничения. Каждый экземпляр приложения хранит полный набор данных справочника. Соответственно при масштабировании сервиса объём данных также увеличивается. Ещё один нюанс: при запуске нового экземпляра Kafka Streams должен полностью восстановить состояние таблицы, прочитав все записи из соответствующего топика. В нашем случае такой подход оказался оптимальным. Если бы речь шла о десятках или сотнях миллионов записей, мы бы рассматривали другие варианты хранения данных.

Что делать, если…

Реальная жизнь редко бывает идеальной. Иногда сделка приходит раньше справочника. Иногда появляется новый инструмент, данные по которому ещё не были загружены. Иногда возникают временные ошибки валидации. В таких случаях мы отправляем сообщение в отдельный DLQ-топик.

Схема выглядит следующим образом:

Сделка → ошибка обогащения → DLQ → повторная публикация → повторная обработка

Для DLQ настроен retry-механизм: через определённый промежуток времени сообщения автоматически возвращаются в основной поток обработки. У нас настроен экспоненциальный бэкоф, то есть сначала мы пробуем направить данные через 10 минут, далее через 1 час, далее через 3 часа и так до 24 часов. К этому моменту справочник обычно уже обновлён, и сделка успешно проходит обогащение. Такой подход позволил отказаться от ручной обработки большинства подобных ситуаций.

Как это работает:

Могут быть и другие проблемы: Kafka Streams работает с одним кластером Kafka. Во время внедрения мы столкнулись ещё с одной особенностью Kafka Streams. Приложение работало в рамках одного Kafka-кластера, тогда как часть наших данных поступала из другого кластера. Нужно было обеспечить доступность справочников в едином контуре обработки.

Мы рассмотрели два варианта:

  • MirrorMaker;

  • собственный сервис репликации.

В итоге остановились на собственном сервисе репликации, который переносит данные из Cluster B в Cluster A. После этого Kafka Streams-приложение работает только с одним кластером и получает все необходимые данные из локальных топиков.


Архитектура

Как видим, наш сервис реплицирует данные из Cluster B в Cluster A, и наш основной сервис обогащения всегда работает с актуальными данными.

Что мы получили?

После перехода на Kafka Streams мы получили:

  • обработку более 3 миллионов сделок в сутки в режиме реального времени;

  • отсутствие обращений к базе данных при обогащении;

  • снижение задержек обработки с часов до секунд;

  • автоматическое восстановление после большинства ошибок через DLQ;

  • горизонтальное масштабирование за счёт увеличения количества инстансов.

Главное изменение оказалось не техническим, а бизнесовым: бизнес перестал ждать окончания торгового дня и получили возможность видеть результаты обработки практически сразу после совершения сделки. Это позволило операционным подразделениям быстрее реагировать на ошибки в справочниках и интеграциях, устранять проблемы с новыми финансовыми инструментами, получать актуальные данные по обработанным сделкам в течение дня, уменьшить количество инцидентов за счёт более раннего выявления некорректных данных.

Kafka Streams хорошо подходит для задач потокового обогащения данных, когда требуется высокая производительность и минимальные задержки. Особенно удачным решением для нас стало использование GlobalKTable, которое позволило полностью исключить обращения к внешним системам во время обработки сделок.

Конечно, у такого подхода есть ограничения: размер справочников должен помещаться в память каждого экземпляра приложения, а работа с несколькими Kafka-кластерами требует дополнительных архитектурных решений. Тем не менее для нашего сценария переход на Kafka Streams позволил превратить ночной batch-процесс в полноценную систему обработки событий в реальном времени без существенного усложнения инфраструктуры.

Комментарии (4)


  1. alexanderfedyukov
    02.09.2026 04:18

    Спасибо за статью! Рассажите пожалуйста как вы решаете проблему изменения справочника "налету"? Т.е. как ведется обработка данных сделок, которые могут содержать ссылки на новую запись в справочнике, эти новые данные в справочный топик опубликованы, но не считались всеми подами с GlobalKTable маппингом, либо считались неодновременно. Виден риск работы с неконсистентным состоянием системы. Как вы ее достигаете в части работы со справочниками?


    1. RSHB_tsyfra Автор
      02.09.2026 04:18

      Спасибо за вопрос! Здесь действительно есть нюанс с согласованностью GlobalKTable. Мы это учитываем на уровне обработки сделки, если при обогащении необходимой записи в GlobalKTable ещё нет, сделка не отбрасывается и не теряется, она отправляется в DLQ и затем повторно обрабатывается. Как показывает практика при retry, GlobalKTable уже содержит данные и сделка успешно обогащается. Таким образом мы не полагаемся на мгновенную синхронизацию наших справочников между подами, а компенсируем возможную рассинхронизацию retry механизмом. Если же запись действительно отсутствует в справочнике, после исчерпания retry попыток, это уже рассматривается как отдельный сценарий ошибки.


  1. nafiskhusnutdin
    02.09.2026 04:18

    Сколько времени заняла миграция сервиса на Kafka Streams? И было ли сложно изучить эту либу членам команды? Потому что стримы работают над обычными консумерами добавляя кучу новой логики и там могут быть много подводных камней


    1. RSHB_tsyfra Автор
      02.09.2026 04:18

      Спасибо за вопрос! На самом деле сама миграция заняла в районе 5 месяцев, мы мигрировали всевозможные рынки на бирже, такие как фондовый, валютный, срочный, помимо этого ещё нам приходилось делать архивацию событий с биржи и реконсиляцию данных, на это тоже ушло время. Наша команда уже имела опыт работы с KafkaStreams, поэтому много времени на изучение этой библиотеки не требовалось.