
Всем привет! Меня зовут Муса. Наша команда занимается витринами данных по товарному учёту.
Каждый день мы доставляем около 10 млрд записей в разных форматах. На этих данных строится различная аналитика, связанная с товарными запасами и движениями экземпляров. Перед нами встала задача: пять раз в день обогащать выгрузку из миллиардов экземплярных остатков дополнительными атрибутами для построения различного рода аналитики. История этих атрибутов уже измерялась десятками миллиардов записей.
Первое решение выглядело просто: положить данные в ClickHouse и сделать JOIN. Но одна выгрузка считалась около 12 часов, а нам нужно было укладываться в десятки минут.
В статье расскажу, про то, как мы смогли сократить время обработки примерно до 33 минут, про ключевой подход при работе с большими объёмами данных, а также попытаюсь донести важность локальности данных на примере реальной задачи.
Коротко про домен
Мы занимаемся товарным учётом и отслеживаем изменения состояния каждой отдельной единицы товара.
Конкретная единица товара называется экземпляром, и у каждого экземпляра есть набор атрибутов:
Уникальный идентификатор экземпляра.
Идентификатор товара.
Локация (На складе A, На складе B, Продан, В пути).
Цена и т. д.
Любое изменение состояния экземпляра — это экземплярное движение.
Погрузили в машину — движение. Продали — движение. Утилизировали — движение.
Мы отслеживаем все экземплярные движения, благодаря этому есть возможность определить состояние экземпляра на конкретный момент.
Самым частым запросом на получение данных является выгрузка экземплярных остатков (экземпляры, которые не проданы и не утилизированы).
Архитектура до решения
В нашей текущей архитектуре есть три ключевых компонента:
1. Мастер‑система по учёту экземплярных движений (exemplar-movement-storage).
Как только изменилось состояние экземпляра (т.е. произошло экземплярное движение), мастер‑система публикует событие в топик Kafka — exemplar_movements.
2. Витрина экземплярных остатков.
Сервис потребляет топик экземплярных движений и прихранивает в базу данных. Пять раз в день запускается процесс расчёта и выгрузки текущих экземплярных остатков в Kafka.
Каждая выгрузка содержит миллиарды экземпляров.
3. Enricher.
Сервис также потребляет топик экземплярных движений и занимается обогащением дополнительными аналитическими атрибутами. К аналитическим атрибутам относятся атрибуты, которые принадлежат внешним системам или для подсчёта которых всё так же нужно сходить во внешние системы. В качестве примера ограничимся одним атрибутом supply_id — это идентификатор поставки, в рамках которой появился экземпляр.
Сервис не хранит в себе результат обогащения, его задача только в самом обогащении и в отправке результата в топик Kafka.

Задача
Выгрузка экземплярных остатков содержит только базовые атрибуты. Этого достаточно для большинства потребителей, но не хватает ряда атрибутов для построения аналитики. Поэтому необходимо обогащать выгрузку экземплярных остатков дополнительными аналитическими атрибутами. Также необходимо уметь производить расчёт не только на текущую, но и на более раннюю дату.
Проектируем верхнеуровневое решение
Итак, у нас есть два источника данных:
Enricher — публикует аналитические атрибуты для экземплярных движений.
Витрина экземплярных остатков — пять раз в день публикует выгрузки актуальных экземплярных остатков, которые содержат миллиарды записей.
Наша задача — связать эти два потока и получить единый обогащённый результат.
Накопление истории
Сервис подписывается на топик обогащённых движений от Enricher'а и сохраняет каждое событие в постоянное хранилище. Зачем? Enricher не хранит данные — он только обогащает и отправляет аналитические атрибуты. Если мы не сохраним результат у себя, то при следующей выгрузке придётся заново ходить во внешние системы, а это сильно скажется на стабильности и скорости выгрузки.
Кроме того, иногда возникают запросы на расчёт остатков за прошлые периоды, например для построения аналитики. Без полной истории это невозможно. Поэтому мы храним все аналитические атрибуты для всех экземплярных движений, а не только для актуальных.
Для этого необходимо разово получить и сохранить все аналитические атрибуты по всем экземплярным движениям, чтобы иметь полную историю по ним
Обогащение выгрузки
Когда витрина публикует очередную выгрузку экземплярных остатков, наш сервис получает её и для каждой строки находит соответствующие аналитические атрибуты из сохранённой истории, после чего публикует обогащённый результат в отдельный топик Kafka.
Если говорить упрощённо, то на втором шаге нам нужно выполнить операцию, аналогичную JOIN по ключу (exemplar_id, operation_id) между двумя наборами данных: выгрузкой экземплярных остатков и историей аналитических атрибутов.
Выбор базы данных
Итак, перед нами две задачи:
Хранить историю обогащённых движений (десятки миллиардов) с постоянным потоком вставки (несколько сотен миллионов записей в день).
Пять раз в день вставлять по несколько миллиардов строк выгрузки и обогащать их через JOIN с исторической таблицей аналитических атрибутов.
База данных должна справляться с такими всплесками без деградации производительности.
С учётом объёмов и характера нагрузки мы сразу смотрели в сторону OLAP‑решений.
У команды был опыт с ClickHouse, поэтому выбрали его, так как он покрывает следующие требования:
Высокая скорость вставки.
Высокая скорость чтения.
Хорошее сжатие.
Нативная поддержка JOIN.
Конфигурация
Поднимаем локальный ClickHouse, вставляем 100к аналитических атрибутов из прода, получаем:
средний размер строки в несжатом формате — 96 байт;
коэффициент сжатия — 5,94.
Также вставляем полную выгрузку из прода, получаем:
средний размер выгрузки в несжатом формате — 37 ГБ;
средний размер выгрузки в сжатом формате — 13 ГБ.
Мы выбрали 3 шарда по две реплики, итого 6 нод.
В качестве ключа шардирования — уникальный идентификатор экземпляра.
Наивное решение
Таблица аналитических атрибутов:
CREATE TABLE analytics_attributes ( exemplar_id Int64, operation_id String, supply_id Int64, // more attributes... created_at DateTime('Etc/UTC') DEFAULT now() ) ENGINE = ReplicatedReplacingMergeTree('/clickhouse/tables/{shard}/{database}/{table}', '{replica}') ORDER BY (exemplar_id, operation_id);
Ключ сортировки — exemplar_id, operation_id, так как джойнить планируем именно по этим полям.
Процесс обогащения выглядит просто:
Приходит выгрузка от витрины (несколько миллиардов строк).
Мы сохраняем её в промежуточную таблицу.
Выполняем JOIN:
SELECT es.exemplar_id, es.operation_id, aa.supply_id, // more attributes FROM exemplar_stocks AS es LEFT JOIN analytics_attributes AS aa ON es.exemplar_id = aa.exemplar_id AND es.operation_id = aa.operation_id
Выглядит идеально: один проход, нет промежуточных таблиц, результат готов на лету, простая система.
Но при запуске обнаружим, что запрос выполняется 12 часов для одной выгрузки, а нам нужно обогащать пять выгрузок в день.
Расследование: как ClickHouse читает данные
ClickHouse — OLAP база данных, где высокая пропускная способность при чтении миллиардов строк гораздо важнее скорости поиска одной записи.
Именно поэтому ClickHouse читает данные с диска не построчно, а крупными пачками.
Крупная пачка называется гранулой, это минимальная единица чтения с диска, и по умолчанию используется адаптивная гранулярность: размер гранулы ограничивается либо количеством строк (по умолчанию 8192), либо объёмом несжатых данных (по умолчанию 10 МБ) — смотря какое ограничение выполнится раньше.
Для более простого восприятия будем рассматривать только ограничение гранулы количеством строк.
При вставки новых данных, ClickHouse не изменяет существующие файлы на диске, вместо этого создаётся новый, отсортированный и сжатый блок данных, который называется парт (не путать с партицией), порядок определяет указанный ORDER BY (…) при создании таблицы.
В качестве первичного индекса используется разреженный индекс (sparse index), что позволяет хранить указатели не на каждую строку, а только на первую строку каждой гранулы. Записи такого индекса называются засечками (marks).
Для поиска конкретной записи ClickHouse считывает индекс, определяет потенциально подходящие гранулы, а затем читает всю гранулу целиком.
Рассмотрим на примере таблицы с аналитическими атрибутами.
Есть три гранулы:
Гранула 0 содержит строки с exemplar_id от 1 до 8192.
Гранула 1 содержит строки с exemplar_id от 8193 до 16385.
Гранула 2 содержит строки с exemplar_id от 16386 до 24578.
Мы хотим получить аналитические атрибуты для трёх экземпляров:
SELECT * FROM analytics_attributes WHERE exemplar_id in (1, 8193, 16386);
Как ClickHouse выполнит этот запрос:
Считывает первичный индекс, определяет, что exemplar_id = 1 лежит в грануле 0, exemplar_id = 8193 в грануле 1, exemplar_id = 16 386 в грануле 2.
Читает все три гранулы целиком с диска (3 × 8192 = 24 576 строк).
В памяти отфильтровывает 24 573 лишние строки, оставляя только 3 нужных.
Ситуация усугубляется, если нужно прочитать много экземпляров в случайном порядке.
В худшем случае каждый переданный exemplar_id окажется в своей грануле, и мы прочитаем 8192 × N строк, где N — количество запрашиваемых экземпляров.
Диагноз: хаотичные чтения с диска
Представим вымышленный пункт выдачи заказов, в котором очень неоптимально хранятся товары на стеллажах.
Вам пришёл заказ из десяти товаров. Вы берёте список и идёте в зону хранения товаров. Но товары разложены хаотично. Вы идёте к первому стеллажу, забираете один товар, идёте к следующему стеллажу, расположенному в конце зала, забираете второй товар, возвращаетесь снова к первому стеллажу. Десять товаров — десять кругов по залу. Большую часть времени вы тратите на перемещение, а не на сборку. Именно это происходит с ClickHouse. Каждая гранула — как стеллаж. Чтобы получить один exemplar_id, ClickHouse читает весь стеллаж целиком (8192 строк), берёт одну коробку, а остальные 8191 игнорируются.

Проблема не в ClickHouse. Проблема в том, что данные физически не подготовлены для такого JOIN. У нас есть две таблицы, которые нужно соединить, но их записи разбросаны по диску независимо друг от друга. ClickHouse — мощный инструмент, но он требует, чтобы данные, которые ищутся вместе, лежали рядом. Иначе он превращается в курьера, который бегает за каждым товаром в другой конец зала.
Виртуальные партиции
Вернёмся к нашей метафоре с пунктом выдачи заказов.
Проблема была в том, что товары разложены хаотично. Чтобы собрать один заказ, приходится бегать по всему залу. А что, если заранее раскладывать товары так, чтобы все позиции из одного заказа лежали на одном стеллаже?
Для этого каждому заказу присваиваем номер стеллажа по простой формуле:
номер стеллажа = номер заказа % количество стеллажей
При поступлении товаров раскладываем их согласно этому правилу. Теперь для сборки заказа достаточно подойти к одному стеллажу и забрать всё сразу.
Переносим идею в ClickHouse
Введём понятие «виртуальная партиция» — это и есть наш стеллаж.
Добавим во все ключевые таблицы поле virtual_partition:
virtual_partition Int32
В качестве формулы определения номера виртуальной партиции будем использовать:
sipHash64(exemplar_id) % 50000
SipHash64 — быстрая хеш‑функция, для любых двух разных exemplar_id получим равномерно распределённые значения. Почему 50 000, будет описано ниже.
Также изменим ORDER BY, добавив в самое начало virtual_partition, именно он задаёт физический порядок строк.
Визуально данные будут храниться следующим образом:

Вот что изменения выше представили нам:
записи про один экземпляр в разных таблицах будут иметь одинаковый virtual_partition (будут находиться на одном стеллаже);
мы можем обрабатывать данные параллельно по виртуальным партициям и быть уверенными, что всё для JOIN'а находится в одной и той же виртуальной партиции.
Почему 50 000
Выбор числа — компромисс между двумя ограничениями:
размер партиции должен быть как минимум больше гранулы (8192 строк);
количество партиций не должно сильно сказываться на эффективности сжатия данных.
Число виртуальных партиций мы подбирали не как магическую константу, а как компромисс между размером одной партиции и их количеством.
Если взять слишком мало виртуальных партиций, обработка каждой партиции будет читать слишком большой объём данных и начнёт упираться в память и I/O.
Если взять слишком много, просядет эффективность сжатия данных.
Целевая архитектура
Перепроектируем обработку выгрузки в четырёхэтапный конвейер.
Высокоуровнево конвейер выглядит так:

Этап 0: Сохранение аналитических атрибутов (фоновый поток от Enricher)
Поток от Enricher идёт постоянно и независимо от выгрузок экземплярных остатков.
Итоговая структура таблицы аналитических атрибутов:
CREATE TABLE analytics_attributes ( exemplar_id Int64, operation_id String, supply_id Int64, virtual_partition Int32 DEFAULT sipHash64(exemplar_id) % 50000, // more attributes... created_at DateTime('Etc/UTC') DEFAULT now() ) ENGINE = ReplicatedReplacingMergeTree('/clickhouse/tables/{shard}/{database}/{table}', '{replica}') ORDER BY (virtual_partition, exemplar_id, operation_id);
Обратите внимание на два изменения:
Добавлено поле virtual_partition, которое вычисляется автоматически при вставке.
ORDER BY начинается с virtual_partition. Теперь данные физически сортируются по стеллажам, а внутри — по экземплярам.
Этап 1: Сохранение экземплярных остатков
Когда Витрина начинает выгрузку, мы получаем данные из Kafka и сохраняем их в таблицу exemplar_stocks.
Партиционируем по calculation_id — ID конкретной выгрузки:
CREATE TABLE exemplar_stocks ( calculation_id UUID, operation_id String, exemplar_id Int64, virtual_partition Int32 DEFAULT sipHash64(exemplar_id) % 50000, more_attributes... ) ENGINE = ReplicatedReplacingMergeTree('/clickhouse/tables/{shard}/{database}/{table}', '{replica}') PRIMARY KEY (virtual_partition) ORDER BY (virtual_partition, exemplar_id, operation_id) PARTITION BY calculation_id;
Почему так:
PARTITION BY calculation_id — каждая выгрузка физически разделена на партиции. Устаревшие выгрузки можно просто удалить с помощью DROP PARTITION, а не чистить через DELETE.
ORDER BY (virtual_partition,...) — внутри выгрузки данные отсортированы по виртуальным партициям. Это даёт локальность для следующего этапа.
virtual_partition вычисляется точно так же, как в таблице аналитических атрибутов.
Для одного и того же exemplar_id — одинаковое значение в обеих таблицах.
Среднее время выполнения: ≈13 минут.
Как только сохранили всю выгрузку экземлярных остатков, приступаем к следующему этапу.
Как устроен пайплайн обогащения
До этого момента обработка была линейной: пришли сообщения из Kafka — записали в ClickHouse.
Теперь начинается самое интересное — обогащение. Оно построено на основе фоновых задач.
Вся работа разбивается на множество независимых фоновых задач. Каждая задача:
имеет тип (подготовка данных, экспорт);
получает на вход идентификатор выгрузки и фиксированный набор виртуальных партиций для обработки;
распределяется между подами, воркерами.
Задачи одного типа одинаковы по логике и различаются только входными данными.
Это даёт два преимущества:
горизонтальное масштабирование (больше подов — быстрее обработка);
устойчивость к сбоям (упал под — перезапускаем только его задачи).
Вернёмся к нашему конвейеру.
Этап 2: Подготовка данных
Это ключевой этап.
Как только выгрузка сохранена, мы создаём 5 000 задач.
Каждой задаче выдаём по 10 виртуальных партиций и calculation_id конкретной выгрузки.
Каждая задача выполняет INSERT... SELECT — но не по всей таблице, а только в пределах своих партиций:
INSERT INTO enriched_exemplar_stocks ( calculation_id, calculation_at, operation_id, exemplar_id, more_attributes..., ) SELECT es.calculation_id, es.calculation_at, es.operation_id, es.exemplar_id, es.more_attributes... FROM ( SELECT exemplar_id, operation_id, supply_id, virtual_partition FROM analytics_attributes WHERE virtual_partition IN {{virtualPartitions:Array(Int32)}} ) AS aa RIGHT JOIN ( SELECT calculation_id, calculation_at, operation_id, exemplar_id, virtual_partition, more_attributes... FROM exemplar_stocks WHERE calculation_id = {{calculationId:UUID}} AND virtual_partition IN {{virtualPartitions:Array(Int32)}} ) AS es ON es.virtual_partition = aa.virtual_partition AND es.exemplar_id = aa.exemplar_id AND es.operation_id = aa.operation_id SETTINGS join_use_nulls = 1;
Почему это работает?
Вспомним наш пункт выдачи. Раньше товары были разбросаны хаотично — чтобы собрать заказ, приходилось бегать по всему залу. Теперь каждый стеллаж (виртуальная партиция) содержит только те товары, которые относятся к нему. Когда приходит задача обогатить выгрузку, мы не бегаем по всему складу — мы подходим к конкретному стеллажу и забираем всё сразу. ClickHouse теперь не нужно прыгать по случайным гранулам. Он знает: все строки с virtual_partition = 42 лежат компактно, в нескольких последовательных гранулах. Он читает только их. Утилизация прочитанных данных в рамках виртуальной партиции близка к 100%.
Таблица обогащённых экземплярных остатков
CREATE TABLE enriched_exemplar_stocks ( calculation_id UUID, calculation_at DateTime64(3, 'Etc/UTC'), operation_id String, exemplar_id Int64, virtual_partition Int32 DEFAULT sipHash64(exemplar_id) % 50000, more_attributes... ) ENGINE = ReplicatedReplacingMergeTree('/clickhouse/tables/{shard}/{database}/{table}', '{replica}') PRIMARY KEY (virtual_partition) ORDER BY (virtual_partition, exemplar_id, operation_id) PARTITION BY calculation_id;
Обратите внимание: virtual_partition сохраняется и здесь — он понадобится для этапа экспорта.
Среднее время выполнения: ≈5 минут.
Как только обработали все виртуальные партиции, приступаем к следующему этапу.
Этап 3: Экспорт
Финальный этап — отправить обогащённые данные потребителям через Kafka.
Создаём 5000 задач экспорта, в каждой — по 10 виртуальных партиций. Каждая задача стримит данные из ClickHouse и отправляет в Kafka:
SELECT calculation_id, operation_id, exemplar_id, more_attributes.... FROM enriched_exemplar_stocks WHERE calculation_id = {calculationId:UUID} AND virtual_partition IN {virtualPartitions:Array(Int32)};
Среднее время: ≈15 минут.
Почему теперь это масштабируется и переживает перезапуски
1. Локальность данных
Главный выигрыш.
Раньше мы просили ClickHouse прыгать по случайным гранулам в поисках одного exemplar_id. Теперь он последовательно читает блоки данных в рамках одной виртуальной партиции.
Утилизация прочитанных строк выросла с ≈0,01% до близкой к 100% в рамках партиции.
2. Параллелизм
50 000 виртуальных партиций дают нам 5 000 независимых задач. Мы можем распределять их между воркерами, и, пока ClickHouse справляется по CPU и I/O, ускорение почти линейное. Упираемся в кластер? Производим решардирование. Упираемся в количество воркеров? Добавляем поды. Архитектура не накладывает ограничений, узким местом становится только железо.
3. Чекпойнты и устойчивость к сбоям
Каждая задача обрабатывает фиксированный набор виртуальных партиций. Если под упал, при перезапуске он перебирает только свои партиции. Весь пайплайн не нужно запускать заново. Это особенно важно для выгрузок с миллиардами строк. Перезапуск с нуля — потерянные часы. Перезапуск нескольких партиций — секунды или минуты.
Результаты
Наивное решение | Целевое решение | |
|---|---|---|
Время одной выгрузки | ≈12 часов | ≈33 минуты |
Утилизация прочитанных строк | ≈0,01% | ≈100% внутри виртуальной партиции |
Параллелизм | Последовательная обработка | 5 000 параллельных задач |
Устойчивость к сбоям | Полный перезапуск | Чекпойнты по партициям |
Количество выгрузок в день | 1 (физически не успевали больше) | 5 (есть запас для 35+ выгрузок) |
Также стоит подметить и минусы целевого решения.
1. Снижение эффективности сжатия данных
Чем больше количество виртуальных партиций, тем хуже коэффициент сжатия данных.
ClickHouse использует колоночное сжатие с алгоритмами (LZ4, ZSTD), которые эффективно работают на повторяющихся последовательностях значений.
Когда мы разбиваем данные на 50 000 виртуальных партиций:
внутри одной виртуальной партиции данные по exemplar_id могут быть сильно разрознены, что снижает эффективность дельта‑кодирования и других методов сжатия;
размер гранул — если партиция становится слишком маленькой (меньше нескольких гранул), сжатие практически не даёт выигрыша.
Поэтому стоит очень тщательно подходить к выбору количества виртуальных партиций.
В нашей задаче для эффективного JOIN'а понадобилось большое количество партиций, но для этапа экспорта можно сильно уменьшить их количество.
2. Необходимость явно указывать virtual_partition в adhoc‑запросах
Для adhoc‑запросов, в которых необходимо получить данные по конкретному exemplar_id, теперь обязательно нужно указывать virtual_partition в условии запроса, иначе ClickHouse не сможет эффективно отфильтровать гранулы.
Вывод
Принцип локальности данных — один из фундаментальных, и он всплывает в самых разных системах, будь то базы данных, распределённые файловые системы или даже кеши.
Если данные, которые обрабатываются вместе, физически лежат рядом, система работает максимально эффективно. Если разбросаны — упираетесь в I/O, и никакая мощь железа не спасёт.
Именно этот принцип мы и применили, и самое главное — подход универсален и не привязан к конкретному инструменту.


