Знакомьтесь, Pulsora: база данных временных рядов на Rust для торговых данных

У нас был поток данных и некуда было его толком деть.

Внутренняя задача была конкретной: хранить и запрашивать большой непрерывный поток рыночных данных — тики, котировки, бары по множеству инструментов — достаточно быстро, чтобы приём данных никогда не становился узким местом, а запросы по диапазону времени возвращались без задержки дыхания. Рыночные данные — это самая чистая и самая жёсткая версия нагрузки временных рядов. Они приходят по порядку, приходят постоянно и не прекращаются никогда. Объём интересен не тем, что он большой один раз. Он интересен тем, что он большой каждую секунду, всегда.

Сначала мы попробовали очевидное: взять существующую базу данных временных рядов. Обычно это правильное решение, и какое-то время всё было нормально. Потом требования к пропускной способности записи поползли вверх, размер на диске рос ещё быстрее, а запросы по диапазону времени, которые мы реально выполняли, начали ощущать вес индекса. В какой-то момент слой хранения делал больше учётной работы на строку, чем полезной работы со строкой.

Поэтому мы создали Pulsora — базу данных временных рядов на Rust, оптимизированную для рыночных данных и упорядоченных по времени наборов данных. Это открытый код на GitHub под лицензией Apache-2.0. Этот пост — введение в то, что это такое, зачем оно существует и как оно на самом деле хранит и запрашивает данные — на основе кода, а не рекламной брошюры.


Почему строить, а не брать готовое

Создание базы данных — не то, за что берутся между делом. Планка, которую нужно было преодолеть: что конкретно нам нужно, чего мы не получали достаточно дёшево?

Три вещи всплывали снова и снова.

Пропускная способность записи, масштабирующаяся по блокам, а не по строкам. Поток тиков — это миллионы маленьких, почти одинаковых строк. Большинство хранилищ общего назначения хотят индексировать каждую из них. Это запись в индексе на строку, а при наших объёмах строк индекс становится доминирующей стоимостью — в амплификации записи, на диске, при компактизации. Мы хотели, чтобы стоимость хранения на строку стремилась к нулю, а индекс масштабировался по числу блоков, которые мы пишем, а не по числу строк.

Сжатие, понимающее данные. Рыночные данные сжимаются превосходно, если относиться к ним как к колонкам известного типа. Метки времени приходят с почти регулярными интервалами. Цены меняются медленно. Объёмы — небольшие целые числа. Универсальное построчное сжатие оставляет почти всё это нетронутым. Мы хотели типозависимое колоночное сжатие — того рода, что описала статья Facebook про Gorilla — встроенное, а не прикрученное сбоку.

Скучный, гибкий по форматам интерфейс. Мы принимаем данные от нескольких разных источников и потребляем в несколько разных инструментов. Нам не нужен был самописный сетевой протокол. Мы хотели делать POST с CSV из скрипта, передавать Apache Arrow из конвейера и читать обратно JSON, Arrow или CSV в зависимости от того, кто спрашивает. HTTP и хорошо известные колоночные форматы, ничего экзотического.

Ни одна из этих вещей сама по себе не нова. Сочетание — для нашей нагрузки, с нашими ограничениями по размеру — оказалось достаточным, чтобы оправдать написание. А раз уж ты всё равно это строишь, написать на Rust ради корректности на этапе компиляции в горячем пути — это самая лёгкая часть. Мы уже писали о том, почему мы строим собственные инструменты; Pulsora ровно в этом ряду.


Что такое Pulsora сегодня

Буду честен насчёт охвата: Pulsora — это 0.1.0. Это движок, который хорошо делает свою работу, а не кластеризованная, реплицируемая, многоарендная платформа с планировщиком запросов. Это одноузловое колоночное хранилище временных рядов на встроенном RocksDB с REST API. Это всё, и это всё — и есть суть.

Вот его форма:

  • Бэкенд хранения: RocksDB, с собственным колоночным слоем поверх.
  • Форматы приёма: CSV, поток Apache Arrow IPC и Protocol Buffers.
  • Форматы вывода запросов: JSON (по умолчанию), поток Arrow IPC, Protobuf и CSV — согласуются через заголовок Accept.
  • Схема: выводится автоматически при первом приёме. Вы не объявляете таблицы; вы делаете POST данных, а Pulsora сама определяет колонки и типы.
  • Интерфейс: HTTP/REST API на Axum поверх Tokio.
  • Долговечность: опциональный журнал упреждающей записи, чтобы буферизованные строки пережили сбой.

Всё это — единственный бинарник с именем pulsora. Есть один флаг командной строки.

# Запуск с настройками по умолчанию
cargo run

# Или указать файл конфигурации
cargo run -- -c pulsora.toml

Этот флаг -c/--config — вся поверхность командной строки. Всё остальное живёт в TOML. Это сделано намеренно — конфигурация и есть API для эксплуатации, а API во время выполнения — это HTTP.


Как данные хранятся на диске

Это та часть, которая нас заботила больше всего, поэтому её и стоит объяснить.

Колоночные блоки, а не строки

Строки приходят, но не хранятся как строки. Pulsora накапливает входящие строки в буфере в памяти, и когда буфер достигает порога по размеру (buffer_size, по умолчанию 1000) или по времени (flush_interval_ms), она группирует их в ColumnBlock — фрагмент строк, хранимый колонка за колонкой, каждая колонка сжата независимо алгоритмом, выбранным под её тип.

Блок — это единица хранения. ColumnBlock примерно такой:

pub struct ColumnBlock {
    pub row_count: usize,                        // строк в этом блоке
    pub columns: HashMap<String, Vec<u8>>,       // сжатые данные колонок
    pub null_bitmaps: HashMap<String, Vec<u8>>,  // отслеживание null по колонкам
}

Хранение по колонкам вместо строк — это то, что заставляет сжатие работать, потому что все значения в колонке имеют один тип и схожий порядок величины. Это и есть вся причина существования колоночных хранилищ, а данные временных рядов — та нагрузка, под которую они и придумывались.

Типозависимое сжатие

Каждый тип колонки получает стратегию сжатия, сделанную под его форму. Согласно документации репозитория:

Тип данных Сжатие Типичный коэффициент (цифры репо) Почему работает
Timestamp Delta-of-delta + varint 5–10x Почти регулярные интервалы
Float XOR (Gorilla) + varfloat 2–5x Медленно меняющиеся значения
Integer Delta + varint 3–8x Последовательные / счётчики
String Словарное кодирование 2–4x Повторяющиеся значения
Boolean Кодирование длин серий (RLE) 10–50x Разреженные данные

Путь меток времени — хрестоматийный случай. Тики приходят с почти постоянными интервалами, поэтому разность второго порядка (delta-of-delta) обычно равна нулю или ничтожно мала, а varint кодирует крошечное число в один байт. Float используют упаковку битов XOR-с-предыдущим-значением из алгоритма Gorilla — когда цена едва двигается, большинство XOR-битов равны нулю и упаковываются в ничто. Строки вроде тикеров повторяются бесконечно, поэтому словарь сопоставляет каждый отдельный символ небольшому целому один раз.

Под всем этим RocksDB применяет второй слой универсального блочного сжатия — lz4 по умолчанию, доступны также none, snappy и zstd. Так вы получаете сначала типозависимое кодирование, а потом проход общего назначения поверх.

Эти коэффициенты — собственные бенчмарк-заявления репозитория, измеренные наборами Criterion в benches/. Я привожу их как цифры проекта, а не как независимый вердикт — запустите cargo bench на собственных данных и доверяйте этому.

Индекс, который не растёт построчно

Это то проектное решение, вокруг которого всё и вертится.

Большинство хранилищ временных рядов держат индекс на строку, чтобы найти любую отдельную строку по времени. Pulsora — нет. Она держит индекс только на уровне блоков — метаданные о каждом блоке, никогда об отдельных строках. Раскладка ключей на диске, согласно документам по архитектуре, выглядит так:

Block index: [table_hash:u32]['B'][min_ts:i64][block_id:u64]
             value: [block_id][min_ts][max_ts][rows][min_id][max_id]   — одна запись на блок
Block data:  [table_hash:u32]['D'][block_id:u64]                       — сжатый колоночный блок
Overrides:   [table_hash:u32]['O'][block_id:u64]                       — позиции «мёртвых» строк

table_hash — это 4-байтовый хеш FNV-1a имени таблицы, используемый как префикс ключа, чтобы таблицы были изолированы без протаскивания строковых имён через каждый ключ. После префикса один байт ('B', 'D', 'O') разделяет индекс, данные и набор переопределений.

Следствие: стоимость хранения масштабируется по числу блоков, а не по числу строк. Запишите миллиард строк — и вы не добавите ни одной записи в индекс сверх блоков, которые их держат. Это и было то требование, которое мы не могли купить дёшево, и это ядро того, зачем Pulsora существует.

Обновления строк обрабатываются без перезаписи блоков. REPLACE отмечает позицию заменённой копии в наборе переопределений старого блока — наборе позиций «мёртвых» строк на блок, объединяемом через merge-оператор RocksDB (объединение множеств) в том же write batch, что и новые данные. Живость — это одно точечное чтение на блок, а не поиск на строку.


Путь записи

Приём — это POST. Схема выводится на первом.

curl -X POST http://localhost:8080/tables/stocks/ingest \
  -H "Content-Type: text/csv" \
  --data-binary @ticks.csv

Строка заголовков этого CSV становится схемой. Pulsora сэмплирует значения, определяет тип каждой колонки (Integer, Float, String, Boolean, Timestamp) и автоматически находит колонку с меткой времени — она понимает RFC 3339, YYYY-MM-DD HH:MM:SS, только дату и Unix-метки времени как в секундах, так и в миллисекундах. С этого момента схема таблицы фиксируется, и входящие данные проверяются на соответствие ей.

Для конвейеров с более высокой пропускной способностью пропустите разбор CSV целиком и передавайте Apache Arrow:

curl -X POST http://localhost:8080/tables/stocks/ingest \
  -H "Content-Type: application/vnd.apache.arrow.stream" \
  --data-binary @ticks.arrow

Arrow приходит со схемой, уже прикреплённой, так что нет шага вывода и нет разбора текста — колонки уже являются колонками. Приём Protobuf работает так же для источников, которые уже им владеют.

Внутри поток такой: разобрать вход, проверить или вывести схему, буферизовать строки в памяти (дедуплицируя на входе, побеждает последняя запись) и — если WAL включён — дописать каждую строку в файл .wal, прежде чем она коснётся памяти, чтобы сбой не потерял буферизованные данные. Когда буфер сбрасывается, строки превращаются в сжатый ColumnBlock, блок и его запись в индексе пишутся в RocksDB одним batch, а WAL усекается.


Путь чтения

Запросы — это GET с диапазоном времени и пагинацией.

# Всё (ограничено лимитом по умолчанию)
curl "http://localhost:8080/tables/stocks/query"

# Временное окно
curl "http://localhost:8080/tables/stocks/query?start=2024-01-01T09:30:00&end=2024-01-01T16:00:00"

# С пагинацией
curl "http://localhost:8080/tables/stocks/query?limit=1000&offset=0"

start и end принимают тот же гибкий набор форматов меток времени, что и приём. limit по умолчанию 1000, максимум не ограничивается; offset пропускает строки для пагинации.

Движок запросов сначала разрешает запрос по индексу блоков: он строит двоичный диапазон ключей из границ времени и итерирует только те блоки, чей [min_ts, max_ts] пересекается с окном. Он собирает блоки-кандидаты, извлекает каждый уникальный блок один раз, распаковывает его, кэширует распакованный результат, чтобы повторный доступ внутри запроса не распаковывал заново, применяет набор переопределений, чтобы пропустить «мёртвые» строки, извлекает нужные строки и сериализует их в том формате, который запросил заголовок Accept.

Хотите результат обратно как Arrow для последующего конвейера вместо JSON? Поменяйте один заголовок:

curl "http://localhost:8080/tables/stocks/query?limit=10000" \
  -H "Accept: application/vnd.apache.arrow.stream" > out.arrow

curl "http://localhost:8080/tables/stocks/query?limit=10000" \
  -H "Accept: text/csv" > out.csv

Тот же запрос, четыре возможных формата вывода, без шага конвертации с вашей стороны. Есть ещё несколько эндпоинтов на чтение — GET /tables для списка таблиц, GET /tables/{name}/schema для выведенной схемы, GET /tables/{name}/count для подсчёта строк, GET /tables/{name}/row/{id} чтобы получить одну строку по id, и GET /health.


Попробуйте

Pulsora собирается на актуальном тулчейне Rust и запускается как единственный бинарник.

# Клонировать и запустить из исходников
git clone https://github.com/muvon/pulsora.git
cd pulsora
cargo run

# Или установить бинарник напрямую
cargo install --git https://github.com/muvon/pulsora.git

# Принять CSV
curl -X POST http://localhost:8080/tables/stocks/ingest \
  -H "Content-Type: text/csv" \
  -d "timestamp,symbol,price,volume
2024-01-01 09:30:00,AAPL,150.00,1000
2024-01-01 09:31:00,AAPL,150.25,1500
2024-01-01 09:32:00,AAPL,149.75,2000"

# Прочитать обратно
curl "http://localhost:8080/tables/stocks/query?limit=10"

# Посмотреть выведенную схему
curl "http://localhost:8080/tables/stocks/schema"

Настройка живёт в pulsora.toml. Параметры, которые стоит знать первыми:

[storage]
data_dir = "./data"
buffer_size = 1000          # строк буферизуется перед записью блока
flush_interval_ms = 1000    # максимальное время ожидания строки в буфере
wal_enabled = true          # буферизованные строки переживают сбой

[ingestion]
batch_size = 10000          # строк на batch записи в базу данных

[performance]
compression = "lz4"         # none | snappy | lz4 | zstd
cache_size_mb = 256         # кэш блоков для чтения

Если вам важно максимальное сжатие в ущерб задержке, поставьте flush_interval_ms = 0 для режима «только пакетами», который удерживает строки, пока буфер не заполнится, давая более крупные и лучше сжатые блоки. Если вам важна задержка чтения на горячих данных, поднимите cache_size_mb. Полный набор параметров задокументирован в doc/CONFIGURATION.md.


Чем она не является и что дальше

Избавлю вас от разочарования узнать это на собственном опыте. Pulsora сегодня одноузловая. Нет кластеризации, нет репликации, нет SQL, нет общего планировщика запросов, нет проталкивания предикатов сверх диапазона времени, нет ограничения частоты запросов. Запросы фильтруют по времени и постранично разбивают; они не делают GROUP BY. Файл doc/ARCHITECTURE.md перечисляет направления, на которые мы смотрим — отсечение колонок и проталкивание предикатов, кэш блоков на диске, многоуровневое горячее/тёплое/холодное хранение и кластеризацию — но это дорожная карта, а не реальность. Я предпочитаю, чтобы вы знали границы, а не натыкались на них.

Чем она является сегодня — это сфокусированный движок, который принимает поток упорядоченных по времени данных, сжимает их жёстко алгоритмами, которые их понимают, и обслуживает запросы по диапазону времени без индекса, который раздувается построчно. Это была наша задача. Это задача, которую Pulsora решает.


Открытый код

Pulsora на GitHub под Apache-2.0. Это Rust — RocksDB для хранения, Axum и Tokio для API, Arrow и Prost для колоночного и protobuf-форматов — и собирается простым cargo build. Наборы бенчмарков лежат в benches/, если хотите измерить производительность приёма и запросов на собственных данных, а не доверять нашим.

Если вы тонете в упорядоченных по времени строках — рыночных тиках, показаниях датчиков, чём угодно, что приходит по порядку и не прекращается — клонируйте, направьте на неё поток и расскажите нам, где она ломается. Это самый быстрый путь к тому, чтобы она стала лучше.

— Don

Pulsora — открытый код на GitHub по адресу github.com/muvon/pulsora под лицензией Apache-2.0. Нашли ошибку или хотите функцию? Откройте issue.