Finam

ziplime deep dive  ·  рыночные данные  ·  async python

// data.history() in ziplime
.. Polars, Parquet & asyncio
.. instead of bcolz

Как форк Zipline переосмыслил хранение рыночных данных — и что это означает для написания торговых алгоритмов сегодня.

Когда Quantopian закрылся в 2020 году, большинство из нас перешли на локальный Zipline. Тот же API, те же концепции — переход казался естественным. Но чем дольше с ним работаешь, тем сложнее игнорировать фундаментальную проблему в том, как Zipline хранит и загружает данные. У этой проблемы было имя: bcolz.


Почему bcolz был проблемой

Zipline хранил рыночные данные в bcolz — колоночном бинарном формате, рассредоточенном по тысячам директорий на файловой системе, по одной на каждый инструмент. Формат был полностью непрозрачным: открыть файл и посмотреть, что внутри, было невозможно. Сама библиотека фактически заброшена и едва работает на Python 3.10+. При инжестировании небольшое несоответствие типов или противоречие в часовых поясах приводило к тому, что бандл тихо падал или давал повреждённый результат — что обнаруживалось только при запуске реального бэктеста.

Zipline + bcolz

  • Заброшенная библиотека, сломана на Python 3.10+
  • Тысячи директорий — по одной на инструмент
  • Непрозрачный бинарный формат, ничего не проверить
  • Ошибки инжестирования сложно диагностировать
  • Синхронный однопоточный пайплайн

ziplime + Parquet

  • Открытый стандарт, читается любым инструментом
  • Один файл на бандл — все инструменты, полная история
  • Predicate pushdown — с диска читается только нужное
  • Ошибки инжестирования понятны и предсказуемы
  • Async-пайплайн на asyncio + aiofiles

Хранилище бандлов: BundleRegistry и метаданные

В ziplime каждый бандл регистрируется в BundleRegistry и описывается JSON-файлом метаданных. Это единый источник истины: по имени бандла система знает, где лежат данные, какой адаптер хранилища использовать, какой диапазон дат охвачен и с какой частотой выполнялось инжестирование.

Вот как выглядит реальный файл метаданных бандла:

// limex_us_minute_data_1774172679.json
      {
        "name": "limex_us_minute_data",
        "version": "1774172679",           // unix timestamp запуска инжестирования
        "bundle_storage_class": "...FileSystemParquetBundleStorage",
        "bundle_storage_data": {
          "base_data_path": "/Users/vyacheslav/.ziplime/data"
        },
        "start_date": "2021-01-01T00:00:00Z",
        "end_date":   "2026-03-03T00:00:00Z",
        "trading_calendar_name": "XNYS",
        "frequency_seconds": 86400.0,      // 1 день
        "data_type": "MARKET_DATA"
      }

Несколько важных моментов. Версия бандла — это unix-timestamp запуска инжестирования: каждое новое инжестирование создаёт новую версию, не удаляя предыдущую автоматически, что позволяет откатиться назад при необходимости. Поле bundle_storage_class содержит полное квалифицированное имя класса адаптера хранилища: ziplime динамически инстанцирует его при загрузке, не захардкоживая. Это значит, что хранилище подключаемо — файловая система, S3, база данных — при условии, что адаптер реализует нужный интерфейс.


data.history(): что изменилось в API

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

Async API

Все методы доступа к данным в ziplime асинхронны. Это означает, что каждый вызов требует await — включая data.history(), data.current() и функции ордеров вроде order_target_percent(). Функция handle_data должна быть объявлена с async def:

handle_data.py
async def handle_data(context, data):
          df    = await data.history(assets=[asset], fields=["close"], bar_count=20)
          price = await data.current(asset, "close")
          await order_target_percent(asset, 0.5)

Это не косметика — это отражение того факта, что весь I/O-пайплайн в ziplime построен на asyncio и aiofiles. Бэктестинг и живая торговля используют одну и ту же кодовую базу, и именно async делает это практичным в масштабе.

Параметры

Параметр Тип Описание
assets list[Asset] Один или несколько инструментов для запроса
fields list[str] Список полей: open, high, low, close, volume и др.
bar_count int Количество баров для возврата
frequency timedelta / Period Частота баров. По умолчанию 1 день. Задаётся через timedelta или строку Period

Одно заметное изменение по сравнению с Zipline: fields теперь принимает список, а не строку. В Zipline для получения нескольких полей требовался отдельный вызов history() на каждое поле или panel-запрос. В ziplime всё запрашивается одним вызовом и возвращается в едином DataFrame.

Возвращаемый тип — pl.DataFrame

Метод возвращает нативный Polars pl.DataFrame, а не объект pandas. Схема фиксированная и предсказуемая:

Returned DataFrame schema

date Date торговая дата или временная метка бара
sid Int64 числовой идентификатор инструмента
close Float64 по одной колонке на каждое запрошенное поле
... отсортировано по ["sid", "date"]

Количество строк всегда равно bar_count × len(assets). Если запросить 20 баров для 3 инструментов, получим 60 строк — все инструменты расположены последовательно, каждый отсортирован по дате внутри своего блока.


Примеры использования

Базовый — один инструмент, одно поле

Простейший случай: один инструмент, одно поле, дневная частота по умолчанию. Результат — DataFrame с колонками date, sid и close. Вызов .to_numpy() на одной колонке даёт обычный одномерный массив, совместимый с любой библиотекой индикаторов на numpy.

один инструмент · одно поле
df = await data.history(assets=[asset], fields=["close"], bar_count=200)
      
      # извлечь как numpy для TA-Lib или аналогов
      prices = df["close"].to_numpy()
      sma    = talib.SMA(prices, timeperiod=50)

Несколько полей за один вызов

Когда индикатору нужно больше одного ценового ряда — ATR, Stochastic, Bollinger Bands — все нужные поля можно получить сразу. Запрос дополнительных колонок не влечёт потерь в производительности: Polars читает из Parquet-файла только те чанки колонок, которые нужны, поэтому добавление high и low в запрос не читает volume или open с диска.

несколько полей · пример ATR
df = await data.history(
          assets=[asset],
          fields=["high", "low", "close"],
          bar_count=14
      )
      
      # каждая колонка извлекается отдельно для TA-Lib
      highs  = df["high"].to_numpy()
      lows   = df["low"].to_numpy()
      closes = df["close"].to_numpy()
      
      atr = talib.ATR(highs, lows, closes, timeperiod=14)

Внутридневные бары через timedelta

Для минутных или часовых стратегий передайте datetime.timedelta в аргумент frequency. Если бандл хранит данные на уровне минут, ziplime агрегирует их до запрошенной частоты на лету — OHLCV-агрегация (первый open, max high, min low, sum volume) происходит внутри ленивого плана Polars до того, как какие-либо данные окажутся в памяти Python.

внутридневные · частота через timedelta
import datetime
      
      # последние 60 минутных баров
      df = await data.history(
          assets=[asset],
          fields=["close"],
          bar_count=60,
          frequency=datetime.timedelta(minutes=1)
      )
      
      # последние 24 часовых бара с объёмом
      df = await data.history(
          assets=[asset],
          fields=["close", "volume"],
          bar_count=24,
          frequency=datetime.timedelta(hours=1)
      )

Недельные бары через Period

Period — это специфичный для ziplime тип, принимающий строки частоты, выровненные по календарю. Используйте его, когда нужны недельные или месячные бары, а не фиксированная временная дельта — неделя не всегда равна 7 × 24 часа с учётом праздников и границ сессий, а Period корректно обрабатывает это, выравниваясь по торговому календарю, указанному в метаданных бандла.

недельные бары · тип Period
from ziplime.constants.period import Period
      
      # 52 недельных бара — год недельных закрытий
      df = await data.history(
          assets=[asset],
          fields=["close"],
          bar_count=52,
          frequency=Period("1w")
      )

Несколько инструментов за один вызов

Все инструменты возвращаются в едином плоском DataFrame. Это существенное отличие от panel-модели Zipline, где каждый инструмент занимал собственную ось. В ziplime инструменты делят одно пространство строк и различаются по колонке sid. Для работы с одним инструментом фильтруйте по sid. Такая структура органично вписывается в групповые операции Polars, если нужно вычислять индикаторы по всем инструментам сразу.

несколько инструментов · плоский DataFrame
df = await data.history(
          assets=[context.aapl, context.msft, context.goog],
          fields=["close", "volume"],
          bar_count=20
      )
      # всего 60 строк (20 × 3), колонки: date, sid, close, volume
      # отсортировано по ["sid", "date"]

      # выделить один инструмент
      aapl_df = df.filter(pl.col("sid") == context.aapl.sid)
      
      # или вычислить что-то по всем инструментам сразу
      last_close = df.group_by("sid").agg(pl.col("close").last())

Что происходит под капотом

Когда поступает вызов data.history(), ziplime открывает Parquet-файл бандла в ленивом режиме через pl.scan_parquet() и сразу добавляет фильтры по запрошенным символам и диапазону дат. Polars проталкивает эти фильтры вниз до уровня row-группы в Parquet-ридере — данные, выходящие за рамки запроса, вообще не читаются с диска.

Если запрошенная частота отличается от хранящейся — например, бандл содержит минутные бары, а алгоритм запрашивает дневные — ресэмплинг происходит внутри того же ленивого плана, до вызова collect(). Вся цепочка от открытия файла до Python DataFrame — это один оптимизированный запрос, а не последовательность шагов Python с промежуточными структурами между ними.

На горячем пути внутри цикла симуляции — где handle_data срабатывает на каждом баре — ziplime поддерживает в памяти индекс позиций строк каждого инструмента в загруженном DataFrame. Получение среза для одного инструмента — это O(1)-поиск, а не сканирование.

«Вся цепочка от открытия файла до Python DataFrame — это один оптимизированный запрос Polars, а не последовательность шагов Python с промежуточными структурами между ними.»

Миграция с Zipline

Логика стратегии переносится с минимальными изменениями: добавьте async/await по всему коду и скорректируйте участки, ожидавшие pandas-объект от data.history(). Единственное, что нельзя мигрировать автоматически — сам бандл: bcolz-бандлы несовместимы с новым форматом и требуют повторного инжестирования из исходных данных. Это разовая затрата, а в плюсе — данные теперь в формате, который можно реально проверять и с которым можно работать осознанно.

Итог

ziplime заменил bcolz на Polars + Parquet и пересобрал весь пайплайн на asyncio. Бандл теперь — это один читаемый файл с явными, самодокументированными метаданными. data.history() возвращает нативный pl.DataFrame, принимает несколько полей за один вызов и поддерживает любую частоту через timedelta или Period. Результат — более честная архитектура: меньше магии, больше предсказуемости, и формат хранения, который будет работать ещё через пять лет.

github.com/Limex-com/ziplime

Open source, активно поддерживается. Исходный код, документация и примеры — всё в одном месте.

Смотреть на GitHub

Автор работал с Quantopian с 2017 по 2020 год, после чего перешёл на локальный Zipline. С 2024 года ziplime является основным фреймворком для продакшн-бэктестинга и живой торговли. Технические детали основаны на изучении исходного кода ziplime и актуальны на момент написания.