Содержание
- Что такое Routine Load в StarRocks
- Lambda и Kappa: куда встает потоковая загрузка
- Routine Load против Kafka Engine в ClickHouse
- Практика: пайплайн из Kafka в StarRocks
- Предусловие: кластер Kafka и доступ из StarRocks
- Готовим окружение Python и продюсер JSON-событий
- Целевая таблица Primary Key для UPSERT
- Создаем задачу CREATE ROUTINE LOAD
- Маппинг вложенных полей через JSONPaths
- Мониторинг через SHOW ROUTINE LOAD
- Что делать при ошибках загрузки
- Заключение
- Референсные ссылки
Аналитика реального времени начинается с того, что данные попадают в хранилище с минимальной задержкой после появления в источнике. Для StarRocks таким мостом из Apache Kafka служит Routine Load, встроенный механизм непрерывной загрузки потока прямо в таблицу. В этой статье соберем рабочий пайплайн: продюсер шлет JSON-события заказов в Kafka, а StarRocks через Routine Load тянет их в таблицу с нативным UPSERT. Разберем архитектурный контекст, сравним подход с Kafka Engine в ClickHouse и покажем мониторинг задачи.
Материал продолжает серию про StarRocks и опирается на ту же базу shop. Kafka берем готовую, кластер 3.9.2 в режиме KRaft. Построение таких пайплайнов требует знания брокеров сообщений, поэтому по ходу мы отмечаем, какие навыки где пригодятся.
Что такое Routine Load в StarRocks
Routine Load это задача, которая живет внутри StarRocks и постоянно вычитывает сообщения из топика Kafka, преобразует их и пишет в целевую таблицу. В отличие от разовой загрузки, задача работает непрерывно и сама следит за прогрессом чтения по партициям. Смещения консьюмера StarRocks хранит в собственных метаданных, а не в Kafka, и коммитит их атомарно вместе с транзакцией загрузки данных.
Из этого вытекает важное свойство. Если нода упала посреди загрузки, после перезапуска задача продолжит с последнего зафиксированного смещения, без потери и без задвоения сообщений. То есть в целевую таблицу каждое сообщение попадает ровно один раз. Управляется задача простыми командами: создать, приостановить, возобновить, остановить.
Проектирование Online-хранилищ данных на StarRocks.
Код курса
STAR
Ближайшая дата курса
7 сентября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800
Lambda и Kappa: куда встает потоковая загрузка
Чтобы понять роль Routine Load, полезно вспомнить два архитектурных паттерна. Lambda разделяет обработку на пакетный слой, который считает точные результаты по историческим данным, и скоростной слой, который дает свежие приблизительные данные с низкой задержкой. Serving-слой объединяет оба результата для пользователя. Минус Lambda это дублирование логики в двух слоях.
Kappa убирает пакетный слой. Единственным источником истины становится журнал событий в Kafka, а любая витрина это результат обработки потока. Если нужно пересчитать данные, поток проигрывают заново с начала топика. Routine Load идеально ложится в Kappa: он превращает поток из Kafka в постоянно обновляемую таблицу StarRocks, которую сразу видят BI-запросы.
Ключевой вопрос потоковых систем это семантика доставки. At-least-once допускает дубли при повторной обработке. Exactly-once гарантирует однократный учет. Routine Load за счет атомарного коммита смещений вместе с данными дает эффективно однократную загрузку в таблицу, а модель Primary Key поверх этого делает повторные события идемпотентными.
Routine Load против Kafka Engine в ClickHouse
ClickHouse решает ту же задачу иначе. Там создается таблица на движке Kafka, которая читает топик, и материализованное представление, которое перекладывает строки из нее в таблицу MergeTree. Смещения при этом управляются через consumer-группу в самой Kafka. Подход рабочий, но у него есть особенности, которые важно знать при выборе.
Первое это реакция на изменение схемы JSON. В StarRocks Routine Load неизвестные поля события просто игнорируются, а пропущенное сопоставленное поле становится null. Продюсер можно расширять новыми полями, не трогая загрузку. В ClickHouse схема жестко задана в таблице и представлении, а одно битое или несоответствующее сообщение способно застопорить консьюмера. Второе это дедупликация. У StarRocks она встроена через Primary Key, у ClickHouse нужен ReplacingMergeTree и ручная логика схлопывания.
Сведем различия в таблицу.
| Аспект | StarRocks Routine Load | ClickHouse Kafka Engine |
|---|---|---|
| Механизм | встроенная задача, тянет из Kafka в таблицу | таблица-движок Kafka плюс материализованное представление |
| Управление смещениями | в метаданных StarRocks, атомарно с загрузкой | через consumer-группу в Kafka |
| Семантика доставки | эффективно однократная в таблицу | at-least-once, возможны дубли |
| Дедупликация | Primary Key UPSERT из коробки | ReplacingMergeTree и ручная логика |
| Изменение схемы JSON | лишние поля игнорируются, пропущенные становятся null | схема жесткая, битое сообщение стопорит консьюмера |
| Управление задачей | SHOW, PAUSE, RESUME, STOP ROUTINE LOAD | операции над таблицей и представлением |
Apache Kafka для инженеров данных
Код курса
DEVKI
Ближайшая дата курса
24 августа, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800
Практика: пайплайн из Kafka в StarRocks
Предусловие: кластер Kafka и доступ из StarRocks
Kafka уже развернута, это кластер 3.9.2 в режиме KRaft на трех брокерах по адресам 10.140.0.91, 10.140.0.92 и 10.140.0.93, порт 9092. Zookeeper тут не нужен, KRaft хранит метаданные сам. Главное сетевое условие: BE-нода StarRocks должна дотягиваться до этих адресов, иначе задача не прочитает ни одного сообщения. Если стенд StarRocks в контейнере, проверьте доступность брокеров изнутри контейнера.
Создадим топик, если его еще нет.
# протестировано для Apache Kafka 3.9.2 (KRaft) kafka-topics.sh --bootstrap-server 10.140.0.91:9092 \ --create --topic orders --partitions 3 --replication-factor 3 # проверка списка топиков kafka-topics.sh --bootstrap-server 10.140.0.91:9092 --list
Готовим окружение Python и продюсер JSON-событий
Продюсер непрерывно шлет события заказов в топик. Обратите внимание на вложенный объект customer, на нем дальше покажем маппинг вложенных полей. Сначала готовим окружение. Библиотека kafka-python ставится в виртуальное окружение, потому что в Ubuntu 24.04 действует правило PEP 668, которое запрещает системному pip ставить пакеты глобально.
# каталог практики и виртуальное окружение mkdir -p ~/article03 && cd ~/article03 python3 -m venv .venv source .venv/bin/activate pip install --upgrade pip pip install kafka-python
Теперь пишем сам продюсер. Файл ~/article03/producer.py.
# протестировано для Python 3.12 и kafka-python 2.0.2
import json, random, time
from datetime import datetime, date, timedelta
from kafka import KafkaProducer
BROKERS = ["10.140.0.91:9092", "10.140.0.92:9092", "10.140.0.93:9092"]
TOPIC = "orders"
STATUSES = ["paid", "shipped", "cancelled"]
EVENTS_PER_SEC = 50
producer = KafkaProducer(
bootstrap_servers=BROKERS,
value_serializer=lambda v: json.dumps(v).encode("utf-8"),
acks="all",
linger_ms=50,
retries=5,
)
def make_event(order_id):
odate = date(2026, 1, 1) + timedelta(days=random.randint(0, 200))
return {
"order_id": order_id,
"customer": {"id": random.randint(1, 20_000_000), "region": random.randint(1, 90)},
"order_date": odate.isoformat(),
"status": random.choice(STATUSES),
"amount": round(random.uniform(1, 5000), 2),
"event_time": datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
}
order_id = 1
delay = 1.0 / EVENTS_PER_SEC
while True:
producer.send(TOPIC, make_event(order_id))
if order_id % 500 == 0:
producer.flush()
order_id += 1
time.sleep(delay)
Запускаем продюсер в отдельном терминале, он пишет примерно 50 событий в секунду. Активируем то же виртуальное окружение и стартуем.
Одно событие выглядит так, это удобно держать перед глазами при настройке JSONPaths.
Целевая таблица Primary Key для UPSERT
Целевую таблицу делаем в модели Primary Key. Это дает нативный UPSERT: повторное событие с тем же order_id обновит строку, а не создаст дубликат. Для витрины реального времени из потока это то, что нужно.
-- протестировано для StarRocks 3.5.0
USE shop;
CREATE TABLE orders_rt (
order_id BIGINT,
customer_id BIGINT,
region_id INT,
order_date DATE,
status VARCHAR(16),
amount DECIMAL(10,2),
event_time DATETIME
) PRIMARY KEY(order_id)
DISTRIBUTED BY HASH(order_id) BUCKETS 16;
Создаем задачу CREATE ROUTINE LOAD
Теперь создаем саму задачу. Она читает топик orders и пишет в orders_rt. Формат сообщений json, а список брокеров и топик указываются в блоке FROM KAFKA. Параметр kafka_default_offsets со значением OFFSET_BEGINNING заставляет прочитать топик с начала.
-- протестировано для StarRocks 3.5.0
CREATE ROUTINE LOAD shop.orders_stream ON orders_rt
COLUMNS(order_id, customer_id, region_id, order_date, status, amount, event_time)
PROPERTIES (
"format" = "json",
"jsonpaths" = "[\"$.order_id\",\"$.customer.id\",\"$.customer.region\",\"$.order_date\",\"$.status\",\"$.amount\",\"$.event_time\"]",
"desired_concurrent_number" = "3",
"max_batch_interval" = "10",
"max_error_number" = "100"
)
FROM KAFKA (
"kafka_broker_list" = "10.140.0.91:9092,10.140.0.92:9092,10.140.0.93:9092",
"kafka_topic" = "orders",
"property.kafka_default_offsets" = "OFFSET_BEGINNING"
);
Проектирование Online-хранилищ данных на StarRocks.
Код курса
STAR
Ближайшая дата курса
7 сентября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800
Маппинг вложенных полей через JSONPaths
Самое интересное это параметр jsonpaths. Он задает, из какого места JSON брать значение для каждой колонки. Порядок путей строго соответствует порядку колонок в блоке COLUMNS. Вложенные поля указываются через точку, поэтому customer.id и customer.region из объекта попадают в плоские колонки customer_id и region_id.
| Колонка таблицы | JSONPath | Значение из события |
|---|---|---|
| order_id | $.order_id | 1 |
| customer_id | $.customer.id | 15782663 |
| region_id | $.customer.region | 36 |
| order_date | $.order_date | 2026-02-11 |
| status | $.status | shipped |
| amount | $.amount | 2862.18 |
| event_time | $.event_time | 2026-07-28 15:43:03 |
Если продюсер добавит в событие новое поле, задача продолжит работать, лишнее поле она просто не заметит. Это и есть устойчивость к дрейфу схемы, о которой шла речь в сравнении с ClickHouse.
Мониторинг через SHOW ROUTINE LOAD
Состояние задачи смотрим командой SHOW ROUTINE LOAD. Поле State должно быть RUNNING. В поле Progress видны текущие смещения по партициям, а поле связанное с отставанием показывает, насколько консьюмер отстал от конца топика. Растущее отставание это сигнал, что загрузка не успевает за потоком.
-- статус задачи целиком SHOW ROUTINE LOAD FOR shop.orders_stream\G -- подзадачи по партициям SHOW ROUTINE LOAD TASK WHERE JobName = "orders_stream"\G -- управление жизненным циклом PAUSE ROUTINE LOAD FOR shop.orders_stream; RESUME ROUTINE LOAD FOR shop.orders_stream; STOP ROUTINE LOAD FOR shop.orders_stream;
Проверяем, что данные реально едут в таблицу.
SELECT count(*) FROM shop.orders_rt; SELECT * FROM shop.orders_rt ORDER BY event_time DESC LIMIT 10;
Что делать при ошибках загрузки
Поток редко бывает идеально чистым. Параметр max_error_number задает, сколько битых сообщений задача стерпит в пределах окна выборки, прежде чем уйти в состояние PAUSED. Если задача встала, причину показывают поля ReasonOfStateChanged и ErrorLogUrls в выводе SHOW ROUTINE LOAD, там же ссылка на лог с отбракованными строками. Типичные причины это несоответствие типа, например строка вместо числа в поле amount, или сообщение, которое вообще не является валидным JSON.
После того как источник поправлен или лимит ошибок поднят, задача возвращается в работу командой RESUME. Смещения при этом не теряются, чтение продолжится с последнего зафиксированного места, поэтому ни одно сообщение не пропадет и не задвоится. Для заведомо мусорных сообщений полезно на стороне продюсера завести отдельный dead-letter топик и не смешивать их с рабочим потоком. Так основная задача остается стабильной, а разбор проблемных событий идет отдельно.
Заключение
Routine Load превращает поток событий из Kafka в живую таблицу StarRocks с минимальной задержкой и однократной загрузкой. За счет модели Primary Key повторные события идемпотентны, а устойчивость к дрейфу схемы JSON снимает боль, знакомую по Kafka Engine в ClickHouse. Такой пайплайн это основа Kappa-архитектуры, где журнал в Kafka остается источником истины, а витрины живут в аналитической СУБД.
Построение подобных конвейеров стоит на трех навыках. Работу с брокерами сообщений мы разбираем на курсе по Apache Kafka для разработчиков. Потоковую аналитику и модель хранения в StarRocks проходим на курсе Проектирование Online-хранилищ данных на StarRocks (код STAR). А тонкости потоковой загрузки в другой аналитической СУБД закрепляем на курсе по ClickHouse (код CLICH). Разобрать эти темы вживую можно на митапе Школы Больших Данных.






