Как вычислить профиль объема (VPVR) на 50 ГБ данных с Dask

Как вычислить профиль объема (VPVR) на 50 ГБ данных с Dask ЛИЧНЫЙ БЮДЖЕТ и ЭКОНОМИЯ
Пошаговое руководство по расчету профиля объема (VPVR) на тиковых данных объемом 50 ГБ с использованием Dask DataFrame без переполнения памяти.
Суть: Расчет профиля объема (VPVR) на 50 ГБ тиковых данных в оперативной памяти невозможен стандартными средствами Pandas из-за OOM. Решение заключается в использовании Dask DataFrame с оптимизированным двухэтапным агрегированием через map_partitions. Это позволяет избежать глобального перемешивания данных (shuffle) и выполнить расчет за минуты на обычном кластере или даже на одной мощной рабочей станции.

Исходный код

При работе с большими объемами тиковых данных ключевая проблема Dask — это фаза Shuffle (перераспределение данных между воркерами). Если выполнять стандартный groupby('price') на 50 ГБ сырых тиков, Dask попытается перегруппировать все строки по цене, что приведет к переполнению памяти. Оптимальный подход — сначала агрегировать данные локально внутри каждой партиции с помощью map_partitions, уменьшая объем данных на 3-4 порядка, и только затем проводить финальное объединение.

Если на одной машине Polars показывает отличные результаты (подробнее в статье Ускорение обработки тиковых данных L2 (Orderbook): Polars против Pandas), то для распределенных вычислений на терабайтных масштабах Dask остается стандартом индустрии.

import dask.dataframe as dd
import pandas as pd
import numpy as np
from dask.distributed import Client

def calculate_local_vpvr(partition: pd.DataFrame, bin_size: float) -> pd.DataFrame:
    """
    Локальная агрегация внутри одной партиции.
    Снижает объем передаваемых по сети данных в тысячи раз.
    """
    if partition.empty:
        return pd.DataFrame(columns=['volume', 'buy_volume', 'sell_volume'], dtype='float64')
    
    # Округляем цены до заданного шага (bin_size)
    partition['price_bin'] = (partition['price'] / bin_size).round() * bin_size
    
    # Разделяем объемы на покупку и продажу
    partition['buy_volume'] = partition['volume'].where(partition['side'] == 1, 0.0)
    partition['sell_volume'] = partition['volume'].where(partition['side'] == -1, 0.0)
    
    # Группируем локально
    local_profile = partition.groupby('price_bin')[['volume', 'buy_volume', 'sell_volume']].sum()
    return local_profile

def compute_global_vpvr(parquet_path: str, bin_size: float = 0.5) -> pd.DataFrame:
    """
    Основная функция для вычисления глобального VPVR без shuffle-эффекта.
    """
    # Загружаем только необходимые колонки для экономии RAM
    df = dd.read_parquet(
        parquet_path, 
        columns=['price', 'volume', 'side']
    )
    
    # Определяем метаданные для map_partitions
    meta = pd.DataFrame(
        columns=['volume', 'buy_volume', 'sell_volume'], 
        dtype='float64'
    )
    meta.index.name = 'price_bin'
    
    # Шаг 1: Локальный VPVR на каждой партиции
    local_profiles = df.map_partitions(
        calculate_local_vpvr,
        bin_size=bin_size,
        meta=meta
    )
    
    # Шаг 2: Финальное объединение локальных профилей
    # Так как размер локальных профилей ничтожно мал, groupby здесь мгновенный
    global_vpvr = local_profiles.groupby(local_profiles.index).sum()
    
    # Запуск вычислений
    result = global_vpvr.compute()
    return result.sort_index()

if __name__ == '__main__':
    # Инициализируем локальный Dask-кластер
    client = Client(n_workers=4, threads_per_worker=2, memory_limit='8GB')
    print(f"Dask Dashboard доступен по адресу: {client.dashboard_link}")
    
    path_to_data = './data/ticks_50gb.parquet'
    vpvr_profile = compute_global_vpvr(path_to_data, bin_size=0.1)
    print(vpvr_profile.head(20))

Разбор параметров

  • parquet_path: Путь к директории с партиционированными Parquet-файлами. Использование Parquet критично, так как он поддерживает колоночное чтение, позволяя игнорировать ненужные фичи.
  • bin_size: Шаг цены для группировки профиля объема (размер тика или кастомный шаг сетки). Чем меньше шаг, тем выше детализация VPVR.
  • side: Направление сделки. В коде принято соглашение: 1 для рыночных покупок (Buy), -1 для рыночных продаж (Sell).
  • map_partitions: Метод Dask, применяющий функцию к каждому Pandas DataFrame (партиции) независимо. Это ключевой элемент оптимизации, предотвращающий OOM.

Как запустить

Для оркестрации подобных ETL-процессов и регулярного пересчета исторических данных отлично подойдет Apache Airflow. Подробнее о настройке таких пайплайнов читайте в статье Apache Airflow для ETL криптосвечей: Оркестрация исторических данных.

Чтобы запустить расчет локально, выполните следующие шаги:

1. Установите необходимые библиотеки: pip install dask[complete] pyarrow pandas numpy.

2. Убедитесь, что ваши тиковые данные сохранены в формате Parquet и разбиты на партиции (например, по дням или часам). Это позволит Dask читать файлы параллельно.

3. Запустите скрипт. Вы сможете отслеживать утилизацию памяти и процессора в реальном времени через Dask Dashboard.

Для мониторинга задержек исполнения в реальном времени обычно используется Интеграция InfluxDB и Grafana для HFT: Мониторинг Latency и Slippage, но для оффлайн-анализа исторических терабайтных датасетов и построения профилей объема связка Dask + Parquet является наиболее эффективным решением.

Оцените статью
FinFluct