Представьте: ваш сайт на 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 даст наибольший эффект.







