Real-Time Analytics

Real-Time Analytics

Real-time analytics (Аналитика в реальном времени) это архитектурный подход к работе с данными, при котором система принимает, обрабатывает и делает доступным для запроса каждое новое событие в течение секунд, а не часов после того, как оно произошло. Термин описывает не отдельный продукт с версией и релизным циклом, а категорию решений, подход реализуют десятки разных стеков технологий, от связки Apache Kafka и ClickHouse до Apache Druid, и у категории нет единого вендора-держателя стандарта. Задача, которую решает real-time analytics, это разрыв между моментом события и моментом, когда на него можно среагировать, например обнаружить мошенническую транзакцию или обновить биржевую котировку раньше, чем она успеет устареть.

 

Что такое Real-Time Analytics

По одному из ведущих обзоров категории, материалу ClickHouse What is Real-Time Analytics? A Complete Guide (обновлено 14 апреля 2026 года), real-time analytics это приём, обработка и выполнение запросов к данным в пределах от миллисекунд до секунд после наступления события. Источник описывает термин как спектр из трёх пересекающихся направлений. Streaming analytics непрерывно считает метрики по потоку, on-demand analytics отвечает на запрос по требованию к уже накопленным свежим данным, а operational analytics встроена прямо в рабочий процесс приложения. Общее у всех трёх направлений одно, ответ приходит, пока событие ещё «свежее» и решение по нему имеет смысл принимать.

Тот же обзор описывает архитектуру шире классической модели «поток, аналитическая БД, дашборд» и прямо включает в контур ML-инференс и governance как обязательные слои. Более ранние материалы по теме, например обзор 2023 года от Imply, коммерческого вендора вокруг Apache Druid, о построении Druid поверх потока, до этого расширения не доходят и ограничиваются связкой потока и аналитической базы. Независимого второго источника на такое расширение определения пока нет, поэтому его стоит воспринимать как обозначившийся тренд, а не устоявшийся стандарт категории.

 

Архитектура и путь данных от события до дашборда

Оба крупных обзора категории сходятся в базовой схеме «поток событий, аналитическая база данных, запросы», расходясь только в детализации. ClickHouse описывает архитектуру четырьмя слоями, Imply в материале про архитектуру Apache Druid обходится двумя укрупнёнными блоками. Ниже разобрана более подробная четырёхслойная модель, потому что она явно называет назначение каждого слоя.

Архитектура real-time аналитики: источник события, потоковый приём Kafka, хранение и индексация в OLAP-базе, обработка SQL-запросов, слой обслуживания дашбордов и алертов

 

Приём и хранение потока (streaming ingestion)

Первый слой это брокер сообщений, например Apache Kafka или Redpanda, который принимает поток событий от источников (клики, транзакции, показания датчиков) и отделяет производителей событий от потребителей. Продюсеру не важно, кто и когда прочитает сообщение, потребитель читает в своём темпе. Второй слой, ingestion & storage, это аналитическая колоночная база данных, которая подписывается на поток брокера и делает данные доступными для запроса сразу по мере поступления, а не по расписанию батч-загрузки. Именно здесь real-time analytics расходится с классическим хранилищем данных (data warehouse), в DWH загрузка это отдельный запланированный шаг ETL, а здесь чтение потока и есть загрузка.

 

Обработка и запрос (query/processing layer)

Третий слой выполняет SQL-запросы к уже загруженным свежим данным, часто опираясь на материализованные представления, которые заранее считают типовые агрегаты и снимают часть нагрузки с движка на момент запроса. Четвёртый слой, presentation, это дашборды, REST API и алертинг, то есть то, с чем реально взаимодействует человек или другая система. Практическое сравнение конкретных технологий этого уровня, ksqlDB против выделенной OLAP-базы, разобрано в статье «Что лучше для аналитики в реальном времени: ksqlDB vs OLAP-база данных?». Глубже познакомиться с построением такого хранилища на практике можно на курсе «Проектирование Online-хранилищ данных на StarRocks», где разбирается именно построение real-time OLAP-хранилища с интеграцией Kafka и CDC.

Проектирование Online-хранилищ данных на StarRocks.

Код курса
STAR
Ближайшая дата курса
26 октября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800

 

Принцип работы, как система удерживает задержку в секундах

Путь одного события состоит из четырёх измеримых шагов. Событие происходит (клик, транзакция, показание датчика), уходит в брокер сообщений, аналитическая база подписана на поток и непрерывно индексирует данные по мере поступления, а SQL-запросы к данным секундной свежести питают дашборды и алерты. Ключевое архитектурное отличие от batch-обработки в том, что в этой цепочке нет этапа «дождаться закрытия окна», то есть паузы перед тем, как накопленный за период набор данных считается завершённым и готовым к агрегации.

На уровне хранилища это означает работу с событиями по одному, а не пачками файлов. Каждое новое сообщение сразу видно последующему запросу. Apache Druid, например, поддерживает exactly-once семантику приёма, чтобы при повторной доставке или сбое соединения событие не терялось и не учитывалось дважды. Архитектурно это реализуется одним из двух классических паттернов. Lambda architecture ведёт batch- и streaming-пайплайны параллельно и сверяет результаты, Kappa architecture обходится единственным потоковым пайплайном без отдельной batch-ветки.

 

Real-time vs batch vs stream processing, где проходит граница

Три термина легко перепутать, потому что все три работают с данными быстро. Разница в том, что именно каждый из них отдаёт как результат. Stream processing, разобранный отдельно в статье про потоковую обработку, это движок непрерывных вычислений над бесконечным потоком, например Apache Flink, который держит состояние и watermark, но сам по себе не отвечает на произвольный запрос пользователя. Real-time analytics это надстройка над таким или похожим потоком, которая добавляет слой хранения и обслуживания произвольных SQL-запросов к свежим данным, а не только заранее заданных агрегатов. Batch, в свою очередь, вообще не претендует на свежесть, он пересчитывает результат заново над уже полностью собранным набором данных.

Критерий Batch processing Stream processing Real-time analytics
Что отдаёт как результат Пересчитанный отчёт над завершённым набором данных Обновляемый агрегат по заранее заданной логике окон и джойнов Ответ на произвольный SQL-запрос к свежим данным по требованию
Типичный инструмент Пакетное ETL-задание, Spark batch Apache Flink, Kafka Streams ClickHouse, Apache Druid, StarRocks поверх потока
Задержка «событие, ответ» Часы-сутки, зависит от расписания запуска Секунды, определяется логикой окна и watermark Секунды, но на уровне произвольного запроса, а не только заданной метрики
Кто формулирует вопрос к данным Заранее известен на этапе разработки конвейера Заранее известен, логика агрегации зашита в код Формулируется на лету аналитиком или дашбордом

Границу удобно проводить не по скорости, а по тому, кто и когда задаёт вопрос к данным. Stream processing отвечает на один заранее известный вопрос непрерывно, real-time analytics держит свежие данные готовыми к любому вопросу, который сформулируют позже.

 

Ограничения и подводные камни

Первое ограничение это стоимость. Постоянно работающий кластер приёма и индексации потока, даже если данные обрабатываются по одному событию, платит за непрерывную работу, а не только за время фактического запроса, в отличие от batch-задания, которое занимает ресурсы только на время запуска. Держать такую инфраструктуру ради отчёта, который нужен раз в сутки и допускает часовую задержку, не окупается.

Второе ограничение это согласованность агрегатов. Пока событие проходит цепочку приёма и индексации, любой запрос, выполненный в этот момент, видит агрегат в промежуточном состоянии, не дождавшемся ещё не пришедших или задержанных сетью записей. Это классическая eventual consistency, результат в итоге верен, но «в итоге» не значит «прямо сейчас», и системе нужен явный механизм вроде watermark, чтобы отличить уже готовый агрегат от собираемого. Подробнее о том, какие критерии стоит проверять при выборе конкретной потоковой базы под такую нагрузку, разобрано в статье «7 критериев выбора потоковой базы данных».

Третье ограничение, когда подход не нужен вовсе. Историческая аналитика, отслеживание трендов за месяцы и регулярная отчётность не выигрывают от секундной задержки, если стейкхолдер готов подождать результат несколько минут, обычная batch-аналитика решает ту же задачу дешевле и с меньшим числом движущихся частей.

 

Сценарии использования и отличительные черты

Условие, при котором real-time analytics оправдан, формулируется просто. У устаревших данных должна быть измеримая цена. Ниже сценарии, где оба разобранных обзора категории сходятся.

  • Обнаружение мошенничества. Решение по транзакции нужно принять до её завершения, а не постфактум по вечернему отчёту.
  • Решения, чувствительные ко времени. Алгоритмический трейдинг и динамическое ценообразование теряют смысл, если котировка или спрос успели измениться.
  • Observability и отладка систем. Инцидент нужно увидеть на панели мониторинга за секунды, а не по логам следующего дня.
  • Рекомендации контента и обнаружение аномалий. Модель должна учитывать действие пользователя прямо сейчас, а не по вчерашней выгрузке.
  • Usage-based pricing и продуктовая аналитика. Тариф или продуктовая метрика считаются по факту использования, а не по батчу в конце периода.
  • IoT и телеметрия. Поток показаний датчиков естественно непрерывен, накопление в батч не отражает физическую природу источника.

Отличительная черта этих сценариев одна. Ответ теряет ценность, если задержался.

 

Практика, замер задержки от события до запроса на Kafka и DuckDB

По традиции весь код, использованный в статье, выкладываем на наш GitHub репозиторий
Аналитика в реальном времени (Real-Time Analytics) from airflow import DAG from airflow.operators.bash import BashOperator from datetime import datetime with DAG( dag_id="spark_submit_demo", start_date=datetime(2025, 1, 1), schedule="@daily", catchup=False ) as dag: run = BashOperator( task_id="run_job", bash_command="spark-submit app.py" ) GitHub code example Аналитика в реальном времени (Real-Time Analytics)

Демо собирает минимальную двухслойную архитектуру из раздела выше своими руками, вместо того чтобы верить цифрам вендоров на слово. Слой приёма это Kafka-брокер, образ apache/kafka:4.3.0 в режиме KRaft, переиспользован со стенда для статьи про event streaming. Слой хранения и запроса это встраиваемый OLAP-движок DuckDB 1.5.5, который играет роль аналитической базы данных прямо в процессе потребителя. Клиент Kafka это confluent-kafka 2.15.0.

Продюсер отправляет 30 событий заказов по одному, с паузами 0.05-0.3 с между ними, чтобы получился неравномерный поток, а не пачка.

# confluent-kafka 2.15.0, Apache Kafka 4.3.0 в KRaft, брокер переиспользован со стенда
# event_streaming, прогнано на стенде 2026-08-30
"""Симулирует непрерывный поток заказов: события идут по одному с паузами, а не пачкой разом.
Запускать вторым: топик создаёт consumer_realtime.py, он же должен подтвердить назначение
партиций в своём выводе («жду события»), прежде чем стартует этот скрипт — иначе первые
события уйдут в топик до появления потребителя и задержка получится не той, что в реальном
потоке."""
import json
import random
import time

from confluent_kafka import Producer

BROKER = "localhost:9092"
TOPIC = "real_time_analytics_events"
EVENT_COUNT = 30
CATEGORIES = ["electronics", "clothing", "grocery", "books"]

random.seed(42)


def main() -> None:
    producer = Producer({"bootstrap.servers": BROKER, "acks": "all"})
    sent = 0

    def on_delivery(err, msg):
        nonlocal sent
        if err is not None:
            print(f"ОШИБКА доставки: {err}")
            return
        sent += 1

    t0 = time.time()
    for i in range(1, EVENT_COUNT + 1):
        event = {
            "event_id": i,
            "event_time": time.time(),  # момент, когда покупка реально произошла
            "category": random.choice(CATEGORIES),
            "amount": round(random.uniform(10, 500), 2),
        }
        producer.produce(TOPIC, value=json.dumps(event).encode(), on_delivery=on_delivery)
        producer.poll(0)  # даёт колбэкам доставки отработать, не дожидаясь их
        time.sleep(random.uniform(0.05, 0.3))  # неравномерный поток вместо пачки
    producer.flush(15)
    elapsed = time.time() - t0

    print(f"отправлено событий: {sent} из {EVENT_COUNT} за {elapsed:.2f} с")


if __name__ == "__main__":
    main()

Потребитель создаёт топик сам при старте, чтобы он точно существовал до подписки, иначе первый опрос падает на UNKNOWN_TOPIC_OR_PART, а после появления топика потребитель вычитывает уже накопившийся хвост одной пачкой, и замер задержки перестаёт отражать реальный поток. На каждое событие он вставляет строку в DuckDB, сразу пересчитывает агрегат по категориям и печатает разницу между временем события и моментом, когда обновлённый агрегат стал доступен запросу.

# confluent-kafka 2.15.0, DuckDB 1.5.5, Apache Kafka 4.3.0 в KRaft на стенде event_streaming,
# прогнано на стенде 2026-08-30
"""Читает поток заказов, сразу пишет каждое событие в DuckDB и пересчитывает агрегат по
категориям. Печатает задержку между моментом события и моментом, когда обновлённый агрегат
стал доступен запросу — это и есть задержка «событие -> инсайт» в real-time аналитике.
Запускать первым: скрипт ждёт назначения партиций, прежде чем producer.py начнёт слать события."""
import json
import os
import time

import duckdb
from confluent_kafka import Consumer
from confluent_kafka.admin import AdminClient, NewTopic

BROKER = "localhost:9092"
TOPIC = "real_time_analytics_events"
DB_PATH = os.path.join(os.path.dirname(__file__), "realtime_analytics.duckdb")
EVENT_COUNT = 30
POLL_TIMEOUT = 1.0
IDLE_LIMIT = 5  # столько пустых опросов подряд после назначения партиций считаем концом потока


def create_topic() -> None:
    """Топик создаётся здесь, а не в producer.py: если подписаться до его создания,
    первый опрос падает на UNKNOWN_TOPIC_OR_PART и часть событий уходит в догоняющую
    пачку вместо равномерного потока — задержка первых событий тогда врёт."""
    admin = AdminClient({"bootstrap.servers": BROKER})
    existing = admin.list_topics(timeout=10).topics
    if TOPIC in existing:
        print(f"топик {TOPIC} уже есть")
        return
    admin.create_topics([NewTopic(TOPIC, num_partitions=1, replication_factor=1)])[TOPIC].result()
    print(f"топик {TOPIC} создан")


def main() -> None:
    create_topic()
    if os.path.exists(DB_PATH):
        os.remove(DB_PATH)  # чистый прогон, старые события не искажают агрегат
    con = duckdb.connect(DB_PATH)
    con.execute("""
        CREATE TABLE events (
            event_id INTEGER,
            event_time DOUBLE,
            category VARCHAR,
            amount DOUBLE
        )
    """)

    consumer = Consumer({
        "bootstrap.servers": BROKER,
        "group.id": "realtime_analytics_dashboard",
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
    })

    state = {"assigned": False}

    def on_assign(_consumer, partitions):
        state["assigned"] = True
        print(f"группа получила партиции {sorted(p.partition for p in partitions)}, жду события", flush=True)

    consumer.subscribe([TOPIC], on_assign=on_assign)

    latencies = []
    idle = 0
    while len(latencies) < EVENT_COUNT and idle 2} ({payload['category']}): задержка {latency:.3f} с", flush=True)

    consumer.close()

    print(f"\nобработано событий: {len(latencies)}")
    if latencies:
        print(f"задержка событие -> доступный агрегат: мин {min(latencies):.3f} с, "
              f"среднее {sum(latencies) / len(latencies):.3f} с, макс {max(latencies):.3f} с")

    print("\nитоговый агрегат по категориям:")
    for category, total, count in con.execute(
        "SELECT category, SUM(amount), COUNT(*) FROM events GROUP BY category ORDER BY category"
    ).fetchall():
        print(f"  {category:9.2f}  заказов {count}")

    con.close()


if __name__ == "__main__":
    main()

 

Построение DWH на ClickHouse

Код курса
CLICH
Ближайшая дата курса
12 октября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800

 

Ключевые строки прогона выглядят так.

событие  1 (electronics): задержка 0.347 с
событие  2 (clothing): задержка 0.224 с
событие  3 (electronics): задержка 0.144 с
событие  4 (electronics): задержка 0.081 с
событие  5 (electronics): задержка 0.013 с
...
событие 30 (clothing): задержка 0.012 с

обработано событий: 30
задержка событие -> доступный агрегат: мин 0.011 с, среднее 0.038 с, макс 0.347 с

итоговый агрегат по категориям:
  books        сумма    756.95  заказов 6
  clothing     сумма   1171.69  заказов 5
  electronics  сумма   3098.90  заказов 11
  grocery      сумма   1918.64  заказов 8

Задержка падает с сотен миллисекунд до единиц-десятков за первые 3-4 события и дальше там и остаётся, минимум 0.011 с, среднее 0.038 с, максимум 0.347 с на все 30 событий. Похоже на разогрев клиента Kafka при выходе на рабочий режим, а не на задержку самой вставки в DuckDB: как только клиент прогревается, задержка каждого следующего события падает и там и остаётся. На таком объёме DuckDB демонстрирует сам механизм, для промышленной нагрузки в сотни тысяч событий в секунду слой хранения обычно заменяют на выделенный кластер, например ClickHouse или Apache Druid.

 

Заключение

Real-time analytics закрывает разрыв между событием и решением по нему, но не бесплатно. Постоянно работающая инфраструктура приёма и индексации стоит дороже периодического batch-задания, а согласованность агрегатов в моменте всегда чуть отстаёт от факта, поэтому подход имеет смысл там, где эта цена окупается, а не по умолчанию для любой аналитики. Демо на Kafka и DuckDB выше показывает механизм в чистом виде, не более сложную архитектуру, а тот же принцип «поток, индексация по событию, запрос к свежим данным», применённый к минимальному набору инструментов.

 

Референсные ссылки

  • ClickHouse, What is Real-Time Analytics? A Complete Guide, обновлено 14 апреля 2026, определение термина, четырёхслойная архитектура, отличие от batch-обработки, ориентиры по задержке и throughput
  • Imply, Real-Time Analytics: Building Blocks and Architecture, опубликовано 18 мая 2023, двухблочная модель на Kafka и Apache Druid, exactly-once семантика, паттерны Lambda и Kappa
  • Apache Kafka Documentation, устройство брокера, топиков, партиций и consumer groups, использованных в демо
  • DuckDB Documentation, встраиваемый OLAP-движок, использованный в демо как слой хранения и запроса