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:
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
Количество строк всегда равно bar_count × len(assets). Если запросить 20 баров для 3 инструментов, получим 60 строк — все инструменты расположены последовательно, каждый отсортирован по дате внутри своего блока.
Polars vs экосистема pandas
Возврат pl.DataFrame означает, что библиотеки, ожидающие numpy-массив или pandas DataFrame — чаще всего TA-Lib — требуют явного преобразования через
.to_numpy() или .to_pandas(). Альтернатива — использовать библиотеки, написанные нативно для Polars, например
polars-talib, которые работают напрямую с Polars-выражениями без промежуточных преобразований и естественным образом выигрывают от векторизованного исполнения.
Примеры использования
Базовый — один инструмент, одно поле
Простейший случай: один инструмент, одно поле, дневная частота по умолчанию. Результат — 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 с диска.
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.
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 корректно обрабатывает это, выравниваясь по торговому календарю, указанному в метаданных бандла.
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, если нужно вычислять индикаторы по всем инструментам сразу.
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)-поиск, а не сканирование.
Миграция с Zipline
Логика стратегии переносится с минимальными изменениями: добавьте async/await по всему коду и скорректируйте участки, ожидавшие pandas-объект от data.history(). Единственное, что нельзя мигрировать автоматически — сам бандл: bcolz-бандлы несовместимы с новым форматом и требуют повторного инжестирования из исходных данных. Это разовая затрата, а в плюсе — данные теперь в формате, который можно реально проверять и с которым можно работать осознанно.
Итог
ziplime заменил bcolz на Polars + Parquet и пересобрал весь пайплайн на asyncio. Бандл теперь — это один читаемый файл с явными, самодокументированными метаданными.
data.history() возвращает нативный
pl.DataFrame, принимает несколько полей за один вызов и поддерживает любую частоту через timedelta или Period.
Результат — более честная архитектура: меньше магии, больше предсказуемости, и формат хранения, который будет работать ещё через пять лет.
Open source, активно поддерживается. Исходный код, документация и примеры — всё в одном месте.
Автор работал с Quantopian с 2017 по 2020 год, после чего перешёл на локальный Zipline. С 2024 года ziplime является основным фреймворком для продакшн-бэктестинга и живой торговли. Технические детали основаны на изучении исходного кода ziplime и актуальны на момент написания.

