Kafka Streams для realtime потокового оброблення даних на сайті

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

Розробка та обслуговування будь-яких видів сайтів:

Інформаційні сайти або веб-програми
Сайти візитки, landing page, корпоративні сайти, онлайн каталоги, квіз, промо-сайти, блоги, ресурси новин, інформаційні портали, форуми, агрегатори
Сайти або веб-програми електронної комерції
Інтернет-магазини, B2B-портали, маркетплейси, онлайн-обмінники, кешбек-сайти, біржі, дропшиппінг-платформи, парсери товарів
Веб-програми для управління бізнес-процесами
CRM-системи, ERP-системи, корпоративні портали, системи управління виробництвом, парсери інформації
Сайти або веб-програми електронних послуг
Дошки оголошень, онлайн-школи, онлайн-кінотеатри, конструктори сайтів, портали надання електронних послуг, відеохостинги, тематичні портали

Це лише деякі з технічних типів сайтів, з якими ми працюємо, і кожен із них може мати свої специфічні особливості та функціональність, а також бути адаптованим під конкретні потреби та цілі клієнта.

Послуги, які ми пропонуємо
Показано 1 з 1Усі 2062 послуг
Kafka Streams для realtime потокового оброблення даних на сайті
Складний
~5 днів

Наші компетенції:

Часті запитання

Останні роботи

  • image_website-b2b-advance_0.webp
    Розробка сайту компанії B2B ADVANCE
    1418
  • image_web-applications_feedme_466_0.webp
    Розробка веб-додатків для компанії FEEDME
    1286
  • image_websites_belfingroup_462_0.webp
    Розробка веб-сайту для компанії БЕЛФІНГРУП
    983
  • image_ecommerce_furnoro_435_0.webp
    Розробка інтернет магазину для компанії FURNORO
    1243
  • image_crm_enviok_479_0.webp
    Розробка веб-додатків для компанії Enviok
    983
  • image_bitrix-bitrix-24-1c_fixper_448_0.webp
    Розробка веб-сайту для компанії ФІКСПЕР
    998

Уявіть: ваш сайт на 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.

Процес роботи

  1. Аналітика: аудит топиків, схем даних та бізнес-вимог.
  2. Проектування: вибір топології, типів вікон, налаштування серіалізації (Avro + Schema Registry).
  3. Реалізація: код топології, state stores, Interactive Queries.
  4. Тестування: TopologyTestDriver, інтеграційні тести з Embedded Kafka.
  5. Деплой: 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 дасть найбільший ефект.