Apache Airflow для ETL криптосвечей: Оркестрация исторических данных

Apache Airflow для ETL криптосвечей: Оркестрация исторических данных ЛИЧНЫЙ БЮДЖЕТ и ЭКОНОМИЯ
Подробное руководство по интеграции Apache Airflow для автоматизации ежедневного ETL процесса загрузки исторических OHLCV свечей с криптобирж, включая код и настройку.
Суть: В этой статье мы разработаем и настроим DAG в Apache Airflow для автоматизированной ежедневной загрузки исторических OHLCV свечей с нескольких криптобирж с использованием библиотеки ccxt, обеспечивая надежную и масштабируемую оркестрацию данных.

В мире высокочастотной торговли и количественного анализа доступ к актуальным и историческим данным криптовалютных бирж является критически важным. Ежедневная загрузка, очистка и хранение этих данных — это рутинная, но сложная задача, требующая надежной автоматизации. Apache Airflow, как мощный инструмент для оркестрации рабочих процессов, идеально подходит для решения этой проблемы, позволяя создавать масштабируемые, отказоустойчивые и наблюдаемые ETL-процессы.

В этом руководстве мы покажем, как интегрировать Apache Airflow для создания DAG (Directed Acyclic Graph), который будет ежедневно загружать исторические OHLCV (Open, High, Low, Close, Volume) свечи с выбранных криптобирж, используя популярную библиотеку ccxt. Мы рассмотрим структуру DAG, реализацию задач на Python и необходимые шаги для запуска.

Исходный код

Ниже представлен полный код DAG для Apache Airflow. Этот DAG состоит из нескольких задач: инициализация, загрузка данных для каждой биржи и символа, а также завершение процесса. Для простоты, данные будут сохраняться в локальную директорию в формате Parquet. В реальных условиях рекомендуется использовать облачное хранилище, такое как S3 или MinIO.


from __future__ import annotations

import pendulum
import os
import pandas as pd
import ccxt

from airflow.models.dag import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.dummy import DummyOperator

# --- Configuration --- #
# Define the exchanges and symbols to fetch
EXCHANGES = ['binance', 'bybit'] # Add more exchanges as needed
SYMBOLS = ['BTC/USDT', 'ETH/USDT'] # Add more symbols as needed
TIMEFRAME = '1h'

# Define the directory to store data
DATA_DIR = os.path.join(os.path.dirname(os.path.abspath(__file__)), 'data')

# Ensure the data directory exists
os.makedirs(DATA_DIR, exist_ok=True)

def _fetch_and_store_candles(exchange_id: str, symbol: str, timeframe: str, data_dir: str, **kwargs):
    """
    Fetches historical OHLCV candles for a given exchange and symbol,
    then stores them as a Parquet file.
    """
    print(f"Fetching {symbol} {timeframe} candles from {exchange_id}...")
    exchange_class = getattr(ccxt, exchange_id)
    exchange = exchange_class({'enableRateLimit': True})

    # Calculate 'since' timestamp for the last 24 hours (or more if needed)
    # For daily ETL, we usually fetch data since the last successful run or a fixed period.
    # Here, we fetch for the last 7 days as an example.
    since_timestamp = exchange.parse8601((pendulum.now('UTC') - pendulum.duration(days=7)).isoformat())

    all_candles = []
    while True:
        try:
            # Fetch candles. 'limit' depends on the exchange, 1000 is common.
            candles = exchange.fetch_ohlcv(symbol, timeframe, since=since_timestamp, limit=1000)
            if not candles:
                break
            all_candles.extend(candles)
            # Update 'since' to the timestamp of the last fetched candle + 1 to avoid duplicates
            since_timestamp = candles[-1][0] + 1
            print(f"Fetched {len(candles)} candles. Total: {len(all_candles)}")
            # Add a small delay to respect rate limits if not handled by enableRateLimit
            # time.sleep(exchange.rateLimit / 1000) 
        except ccxt.NetworkError as e:
            print(f"Network error fetching {symbol} from {exchange_id}: {e}. Retrying...")
            # Implement retry logic or raise AirflowSkipException
            raise e # Airflow will handle retries based on task settings
        except ccxt.ExchangeError as e:
            print(f"Exchange error fetching {symbol} from {exchange_id}: {e}. Skipping symbol.")
            return # Skip this symbol if there's an exchange-specific error
        except Exception as e:
            print(f"An unexpected error occurred: {e}. Skipping symbol.")
            return

    if not all_candles:
        print(f"No candles fetched for {symbol} from {exchange_id}.")
        return

    df = pd.DataFrame(
        all_candles,
        columns=['timestamp', 'open', 'high', 'low', 'close', 'volume']
    )
    df['timestamp'] = pd.to_datetime(df['timestamp'], unit='ms', utc=True)
    df = df.set_index('timestamp').sort_index()

    # Remove duplicates based on index (timestamp)
    df = df[~df.index.duplicated(keep='last')]

    # Optional: Data cleaning and validation steps
    # For more advanced cleaning, refer to: 
    # https://finfluct.com/data-analysis/ochistka-ohlcv-dannyh-v-pandas-udalenie-vybrosov-i-nan/
    # For performance, consider Polars: 
    # https://finfluct.com/data-analysis/uskorenie-obrabotki-tikovyh-dannyh-l2-orderbook-polars-protiv-pandas/

    # Define output path
    output_dir = os.path.join(data_dir, exchange_id, symbol.replace('/', '_'))
    os.makedirs(output_dir, exist_ok=True)
    file_path = os.path.join(output_dir, f"{timeframe}_candles.parquet")

    # Append to existing file or create new one
    if os.path.exists(file_path):
        existing_df = pd.read_parquet(file_path)
        combined_df = pd.concat([existing_df, df]).drop_duplicates(subset=['open', 'high', 'low', 'close', 'volume'], keep='last')
        combined_df = combined_df[~combined_df.index.duplicated(keep='last')].sort_index()
        combined_df.to_parquet(file_path, index=True)
        print(f"Appended {len(df)} new candles to {file_path}. Total rows: {len(combined_df)}")
    else:
        df.to_parquet(file_path, index=True)
        print(f"Saved {len(df)} candles to {file_path}")


with DAG(
    dag_id='crypto_ohlcv_etl',
    start_date=pendulum.datetime(2023, 1, 1, tz="UTC"),
    schedule='@daily', # Run once a day
    catchup=False,
    tags=['crypto', 'etl', 'data_ingestion'],
    default_args={
        'owner': 'airflow',
        'depends_on_past': False,
        'email_on_failure': False,
        'email_on_retry': False,
        'retries': 3,
        'retry_delay': pendulum.duration(minutes=5),
    },
) as dag:
    start_etl = DummyOperator(task_id='start_etl')

    # Create a task for each exchange and symbol combination
    fetch_tasks = []
    for exchange_id in EXCHANGES:
        for symbol in SYMBOLS:
            task_id = f"fetch_and_store_{exchange_id}_{symbol.replace('/', '_')}"
            fetch_task = PythonOperator(
                task_id=task_id,
                python_callable=_fetch_and_store_candles,
                op_kwargs={
                    'exchange_id': exchange_id,
                    'symbol': symbol,
                    'timeframe': TIMEFRAME,
                    'data_dir': DATA_DIR,
                },
            )
            fetch_tasks.append(fetch_task)

    end_etl = DummyOperator(task_id='end_etl')

    # Define task dependencies
    start_etl >> fetch_tasks >> end_etl

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

  • dag_id: Уникальный идентификатор DAG. Используется Airflow для отслеживания и управления рабочим процессом.
  • start_date: Дата, с которой Airflow начнет планировать запуски DAG. Задана с помощью pendulum.datetime для удобства работы с временными зонами.
  • schedule: Интервал, с которым DAG будет запускаться. '@daily' означает один раз в день. Можно использовать cron-выражения или None для ручного запуска.
  • catchup: Если False, Airflow не будет запускать пропущенные DAG-запуски между start_date и текущей датой. Рекомендуется для ETL, где нужны только актуальные данные.
  • tags: Список тегов для категоризации DAG в пользовательском интерфейсе Airflow.
  • default_args: Словарь аргументов по умолчанию, которые будут применены ко всем операторам в DAG. Включает retries (количество попыток при сбое) и retry_delay (задержка между попытками).
  • EXCHANGES: Список строковых идентификаторов криптобирж (например, 'binance', 'bybit'), с которых будут загружаться данные. Используется библиотекой ccxt.
  • SYMBOLS: Список торговых пар (например, 'BTC/USDT', 'ETH/USDT'), для которых будут загружаться свечи.
  • TIMEFRAME: Временной интервал свечей (например, '1h', '1d').
  • DATA_DIR: Путь к директории, где будут сохраняться загруженные данные. В примере это поддиректория data рядом с файлом DAG.
  • _fetch_and_store_candles: Python-функция, выполняющая основную логику загрузки данных. Она использует ccxt для подключения к бирже, получения OHLCV данных и сохранения их в Parquet-файл. Включает базовую обработку ошибок и логику дозаписи данных.
  • PythonOperator: Оператор Airflow, который позволяет выполнять произвольный Python-код. В нашем случае, он вызывает функцию _fetch_and_store_candles для каждой комбинации биржи и символа.
  • DummyOperator: Простой оператор, который ничего не делает, но полезен для обозначения начала и конца DAG или для группировки задач.

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

Для запуска этого DAG в вашей среде Apache Airflow выполните следующие шаги:

  1. Установите зависимости: Убедитесь, что у вас установлены необходимые библиотеки. Вы можете создать файл requirements.txt:

    
    apache-airflow
    pandas
    ccxt
    pyarrow # For Parquet support
    

    Затем установите их:

    
    pip install -r requirements.txt
    
  2. Разместите файл DAG: Сохраните приведенный выше код в файл с расширением .py (например, crypto_ohlcv_etl.py) и поместите его в директорию DAGs вашего Airflow. Обычно это директория, указанная в конфигурации Airflow как dags_folder.

  3. Проверьте директорию данных: Убедитесь, что директория data будет создана рядом с файлом DAG или измените путь DATA_DIR на подходящий для вашей среды. Убедитесь, что у пользователя Airflow есть права на запись в эту директорию.

  4. Включите DAG в UI Airflow: Откройте веб-интерфейс Airflow. Вы должны увидеть новый DAG с именем crypto_ohlcv_etl. Переключите его статус на ‘On’ (включено).

  5. Запустите DAG: Вы можете дождаться запланированного запуска (@daily) или запустить его вручную, нажав кнопку ‘Trigger DAG’ в UI Airflow. Следите за логами задач, чтобы убедиться в успешном выполнении.

  6. Проверьте данные: После успешного выполнения задач, вы найдете загруженные Parquet-файлы в директории data, структурированные по биржам и символам (например, data/binance/BTC_USDT/1h_candles.parquet).

Этот DAG предоставляет надежную основу для вашего ETL-процесса. Для дальнейшего улучшения вы можете добавить задачи для валидации данных, уведомлений об ошибках, а также интеграцию с более продвинутыми хранилищами данных. После успешной загрузки и очистки данных, вы сможете использовать их для различных аналитических задач, например, для бэктестинга стратегий Pairs Trading.

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