Представьте: ваш сайт на React генерирует миллионы событий в день — клики, просмотры, покупки. Вы хотите видеть DAU в реальном времени, обогащать заказы данными пользователей и детектить аномалии. Но batch-обработка на Hadoop не успевает: задержка — часы. Мы — команда инженеров с 5+ годами опыта в Kafka Streams, реализовали более 10 проектов потоковой обработки для high-load сайтов. Наше решение — библиотека Kafka Streams, которая работает внутри вашего JVM-приложения без отдельной инфраструктуры. В отличие от Apache Flink или Spark Streaming, здесь не нужно разворачивать кластер — только зависимость в pom.xml. Гарантируем снижение задержки до 10–50 миллисекунд против минут у batch-решений. При нагрузке 50 000 событий в секунду затраты на инфраструктуру снижаются вдвое за счёт отказа от отдельного кластера.
Архитектурная картина
Kafka Streams читает топики, трансформирует, агрегирует, джойнит данные и пишет результат обратно в Kafka или во внешние системы через Kafka Connect. Состояние хранится локально в RocksDB и реплицируется в changelog-топики — это даёт отказоустойчивость без внешней базы. Типичные задачи: агрегация событий пользователей (DAU, воронки), обогащение потока заказов данными из справочников, fraud detection, материализованные представления из event-sourced данных. Мы спроектируем топологию под вашу нагрузку — от тысяч до 300 000 событий в секунду.
Построение топологии обработки
Базовая топология
StreamsBuilder builder = new StreamsBuilder(); KStream<String, UserEvent> events = builder.stream( "user-events", Consumed.with(Serdes.String(), userEventSerde) ); // Фильтрация + трансформация KStream<String, PageView> pageViews = events .filter((userId, event) -> event.getType().equals("PAGE_VIEW")) .mapValues(event -> PageView.from(event)); // Ветвление потока Map<String, KStream<String, UserEvent>> branches = events.split(Named.as("branch-")) .branch((k, v) -> v.getType().equals("PURCHASE"), Branched.as("purchases")) .branch((k, v) -> v.getType().equals("CLICK"), Branched.as("clicks")) .defaultBranch(Branched.as("other")); branches.get("branch-purchases").to("purchase-events"); Агрегации с оконными функциями
Задача — считать количество просмотров страниц по пользователям в скользящем 5-минутном окне:
KTable<Windowed<String>, Long> pageViewCounts = pageViews .groupByKey(Grouped.with(Serdes.String(), pageViewSerde)) .windowedBy( SlidingWindows.ofTimeDifferenceAndGrace( Duration.ofMinutes(5), Duration.ofSeconds(30) // grace period для поздних событий ) ) .count(Materialized.<String, Long, WindowStore<Bytes, byte[]>>as("page-view-counts") .withValueSerde(Serdes.Long()) ); // Publish результатов pageViewCounts.toStream() .map((windowedKey, count) -> KeyValue.pair( windowedKey.key(), new PageViewStat(windowedKey.key(), windowedKey.window().start(), count) )) .to("page-view-stats", Produced.with(Serdes.String(), pageViewStatSerde)); Выбор типа окна зависит от бизнес-логики. Сравнение окон в таблице:
| Тип окна | Поведение | Use case |
|---|---|---|
| Tumbling Windows | Фиксированные непересекающиеся интервалы | Подсчёт событий за каждую минуту |
| Hopping Windows | Пересекающиеся интервалы с фиксированным шагом | Скользящее среднее за 5 минут с обновлением каждую минуту |
| Sliding Windows | Окно сдвигается по каждому событию | Обновление статистики в реальном времени при каждом событии |
| Session Windows | Группировка по периодам активности | Анализ сессий пользователя |
KTable и материализованные представления
KTable — changelog-stream, где каждый новый record с тем же ключом перезаписывает предыдущий. Используется для справочных данных:
KTable<String, UserProfile> userProfiles = builder.table( "user-profiles", Materialized.as("user-profiles-store") ); KStream<String, EnrichedEvent> enriched = events.join( userProfiles, (event, profile) -> EnrichedEvent.builder() .event(event) .userName(profile.getName()) .userSegment(profile.getSegment()) .build() ); Конфигурация приложения
Properties props = new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "site-analytics-processor"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-1:9092,kafka-2:9092,kafka-3:9092"); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.StringSerde.class); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.StringSerde.class); // Производительность props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 4); props.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 10 * 1024 * 1024L); // 10MB props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000); // Обработка ошибок props.put(StreamsConfig.DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG, LogAndContinueExceptionHandler.class); // RocksDB state store props.put(StreamsConfig.STATE_DIR_CONFIG, "/var/lib/kafka-streams"); KafkaStreams streams = new KafkaStreams(builder.build(), props); streams.start(); // Graceful shutdown Runtime.getRuntime().addShutdownHook(new Thread(streams::close)); Interactive Queries — чтение состояния без Kafka
Позволяет читать state store напрямую
ReadOnlyKeyValueStore<String, Long> store = streams.store( StoreQueryParameters.fromNameAndType( "page-view-counts", QueryableStoreTypes.keyValueStore() ) ); Long count = store.get(userId); // Для windowed store ReadOnlyWindowStore<String, Long> windowStore = streams.store( StoreQueryParameters.fromNameAndType( "page-view-counts-windowed", QueryableStoreTypes.windowStore() ) ); WindowStoreIterator<Long> iterator = windowStore.fetch( userId, Instant.now().minus(Duration.ofMinutes(5)), Instant.now() ); Как Kafka Streams сравнивается с Flink и Spark?
| Параметр | Kafka Streams | Flink | Spark Streaming |
|---|---|---|---|
| Инфраструктура | Только JVM-зависимость | Отдельный кластер | Отдельный кластер |
| Задержка | < 50 мс | < 100 мс | > 1 с |
| Масштабирование | Потоки + инстансы | TaskManager | Executor |
| Состояние | RocksDB + changelog | RocksDB/Flink State | Spark State |
| Стоимость (100k/с) | ~$500/мес на инстанс | ~$2000/мес | ~$1500/мес |
Почему стоит выбрать Kafka Streams для веб-аналитики?
Flink требует развёртывания кластера (JobManager + TaskManagers), что увеличивает стоимость и сложность. Kafka Streams работает как обычная библиотека — вы запускаете её вместе с вашим API на том же сервере. Для сайта с нагрузкой до 100 000 событий в секунду Streams справляется без отдельной инфраструктуры. При росте нагрузки масштабируйте увеличением потоков (NUM_STREAM_THREADS) или добавлением инстансов — состояние автоматически балансируется через Kafka.
Как гарантировать, что данные не потеряются при сбое?
Используем exactly-once семантику (processing.guarantee=exactly_once_v2) и state store с changelog-топиками. При перезапуске приложение восстанавливает состояние из Kafka. Настраиваем grace period для поздних событий и Dead Letter Queue для ошибочных записей. В production обязательно мониторим метрики: process-rate, commit-latency, rocksdb-block-cache-hit-ratio — через JMX или Prometheus.
Процесс работы
- Аналитика: аудит топиков, схем данных и бизнес-требований.
- Проектирование: выбор топологии, типов окон, настройка сериализации (Avro + Schema Registry).
- Реализация: код топологии, state stores, Interactive Queries.
- Тестирование: TopologyTestDriver, интеграционные тесты с Embedded Kafka.
- Деплой: Docker-образ с JMX Exporter, CI/CD пайплайн.
Что входит в работу
- Архитектурная схема пайплайна
- Настроенные топики и Schema Registry
- Код топологии с unit-тестами
- Документация по развёртыванию и мониторингу
- Обучение команды (workshop на 1 день)
Сроки
Базовая топология — от 3 до 4 дней. Пайплайн с агрегациями и Interactive Queries — от 6 до 9 дней. Полноценное production-решение с мониторингом, DLQ и CI/CD — от 2 до 3 недель. Свяжитесь с нами для оценки вашего проекта — мы подберём оптимальное решение.
Типичные ошибки при внедрении
- Неправильный выбор окон: для real-time лучше Sliding, не Tumbling.
- Отсутствие grace period — поздние события теряются.
- Слишком большой cache (CACHE_MAX_BYTES_BUFFERING) — увеличивает задержку коммита.
- Игнорирование сериализации: Avro со Schema Registry обязателен.
- Отсутствие мониторинга RocksDB — кеш-хиты падают при переполнении памяти.
Получите консультацию: наши инженеры помогут спроектировать пайплайн под вашу нагрузку. Закажите анализ текущей архитектуры — мы покажем, где Kafka Streams даст наибольший эффект.







