Содержание
- Что такое Real-Time Analytics
- Архитектура и путь данных от события до дашборда
- Приём и хранение потока (streaming ingestion)
- Обработка и запрос (query/processing layer)
- Принцип работы, как система удерживает задержку в секундах
- Real-time vs batch vs stream processing, где проходит граница
- Ограничения и подводные камни
- Сценарии использования и отличительные черты
- Практика, замер задержки от события до запроса на Kafka и DuckDB
- Заключение
- Референсные ссылки
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 обходится двумя укрупнёнными блоками. Ниже разобрана более подробная четырёхслойная модель, потому что она явно называет назначение каждого слоя.
Приём и хранение потока (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 репозиторий
Демо собирает минимальную двухслойную архитектуру из раздела выше своими руками, вместо того чтобы верить цифрам вендоров на слово. Слой приёма это 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-движок, использованный в демо как слой хранения и запроса

