Привет, Хабр! Я Артём Борисов, 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-процесс в полноценную систему обработки событий в реальном времени без существенного усложнения инфраструктуры.


