Уявіть: ваш сайт на React генерує мільйони подій на день — кліки, перегляди, покупки. Ви хочете бачити DAU в реальному часі, збагачувати замовлення даними користувачів та виявляти аномалії. Але batch-обробка на Hadoop не встигає: затримка — години. Ми — команда інженерів з 10+ роками досвіду в Kafka Streams, реалізували більше 40 проектів потокової обробки для 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 дасть найбільший ефект.







