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 выполните следующие шаги:
-
Установите зависимости: Убедитесь, что у вас установлены необходимые библиотеки. Вы можете создать файл
requirements.txt:apache-airflow pandas ccxt pyarrow # For Parquet supportЗатем установите их:
pip install -r requirements.txt -
Разместите файл DAG: Сохраните приведенный выше код в файл с расширением
.py(например,crypto_ohlcv_etl.py) и поместите его в директорию DAGs вашего Airflow. Обычно это директория, указанная в конфигурации Airflow какdags_folder. -
Проверьте директорию данных: Убедитесь, что директория
dataбудет создана рядом с файлом DAG или измените путьDATA_DIRна подходящий для вашей среды. Убедитесь, что у пользователя Airflow есть права на запись в эту директорию. -
Включите DAG в UI Airflow: Откройте веб-интерфейс Airflow. Вы должны увидеть новый DAG с именем
crypto_ohlcv_etl. Переключите его статус на ‘On’ (включено). -
Запустите DAG: Вы можете дождаться запланированного запуска (
@daily) или запустить его вручную, нажав кнопку ‘Trigger DAG’ в UI Airflow. Следите за логами задач, чтобы убедиться в успешном выполнении. -
Проверьте данные: После успешного выполнения задач, вы найдете загруженные Parquet-файлы в директории
data, структурированные по биржам и символам (например,data/binance/BTC_USDT/1h_candles.parquet).
Этот DAG предоставляет надежную основу для вашего ETL-процесса. Для дальнейшего улучшения вы можете добавить задачи для валидации данных, уведомлений об ошибках, а также интеграцию с более продвинутыми хранилищами данных. После успешной загрузки и очистки данных, вы сможете использовать их для различных аналитических задач, например, для бэктестинга стратегий Pairs Trading.




