Исторические ценовые данные

Обзор

Это руководство поможет вам получить доступ к файлам исторических ценовых данных из S3-бакета cfg-public-proper-wallaby (eu-west-1), настроенного с политикой Requester Pays, обеспечивая оптимальную производительность загрузки и стабильность системы.

Предварительные требования

  • Установленный и настроенный AWS CLI
  • Действующие учётные данные AWS с соответствующими правами доступа
  • Доступ к региону eu-west-1
  • Понимание модели биллинга Requester Pays

Шаги настройки

1. Установка учётных данных AWS и региона

export AWS_ACCESS_KEY_ID="your-access-key"
export AWS_SECRET_ACCESS_KEY="your-secret-key"
export AWS_DEFAULT_REGION="eu-west-1"

2. Проверка доступа к бакету

aws s3 ls s3://cfg-public-proper-wallaby/ \
  --region eu-west-1 \
  --request-payer requester

3. Определение доступных инструментов

Перед загрузкой перечислите все валютные пары (инструменты), доступные в бакете. Используйте --delimiter / (или обычную команду aws s3 ls), чтобы вывести только префиксы верхнего уровня — избегайте рекурсивного листинга, который просканировал бы каждый файл только ради ~20-30 имён папок.

Быстрый список (проще всего):

aws s3 ls s3://cfg-public-proper-wallaby/ \
  --region eu-west-1 \
  --request-payer requester

Чистый список (только названия, отсортированные, сохранённые в файл):

aws s3api list-objects-v2 \
  --bucket cfg-public-proper-wallaby \
  --delimiter "/" \
  --region eu-west-1 \
  --request-payer requester \
  --query "CommonPrefixes[].Prefix" \
  --output text | tr '\t' '\n' | sed 's|/$||' | sort > instruments.txt

echo "Всего инструментов: $(wc -l < instruments.txt)"
cat instruments.txt

Python (Boto3):

import boto3
from botocore.config import Config

config = Config(region_name='eu-west-1')
s3_client = boto3.client('s3', config=config)

def list_instruments(bucket):
    """Получить все префиксы инструментов верхнего уровня (валютные пары)"""
    instruments = []
    paginator = s3_client.get_paginator('list_objects_v2')

    page_iterator = paginator.paginate(
        Bucket=bucket,
        Delimiter='/',
        RequestPayer='requester'
    )

    for page in page_iterator:
        for prefix in page.get('CommonPrefixes', []):
            instrument = prefix['Prefix'].rstrip('/')
            instruments.append(instrument)

    return sorted(instruments)

if __name__ == "__main__":
    BUCKET = 'cfg-public-proper-wallaby'
    instruments = list_instruments(BUCKET)

    print(f"Всего найдено инструментов: {len(instruments)}\n")
    for inst in instruments:
        print(inst)

    with open('instruments.txt', 'w') as f:
        f.write('\n'.join(instruments))

Использование:

python3 list_instruments.py

⚠️ Примечание о производительности и стоимости: Использование --delimiter / (или Delimiter='/' в Boto3) обязательно — это вызывает облегчённый ответ CommonPrefixes вместо перечисления каждого файла. Это стоит всего нескольких дешёвых LIST-запросов и занимает секунды, в отличие от медленного и более дорогого полного сканирования бакета без разделителя.

4. Базовая команда доступа к S3 с Requester Pays

# Загрузка одного файла
aws s3 cp s3://cfg-public-proper-wallaby/EURUSD/file.json . \
  --region eu-west-1 \
  --request-payer requester

5. Пакетная загрузка с оптимизацией производительности

# Загрузка всей истории для конкретной валютной пары
aws s3 sync s3://cfg-public-proper-wallaby/EURUSD/ ./EURUSD-data/ \
  --region eu-west-1 \
  --request-payer requester \
  --no-progress \
  --max-concurrent-requests 20 \
  --max-bandwidth 100MB/s

Параметры оптимизации производительности

Параметр Значение Преимущество
--max-concurrent-requests 20-30 Увеличивает число параллельных загрузок
--max-bandwidth 100MB/s Предотвращает сетевые узкие места
--no-progress - Снижает нагрузку на ввод-вывод
--region eu-west-1 Снижает задержку внутри региона ЕС

6. Продвинутый скрипт загрузки (скорость и стабильность)

#!/bin/bash

BUCKET_NAME="cfg-public-proper-wallaby"
REGION="eu-west-1"
CURRENCY_PAIR="${1:-EURUSD}"  # По умолчанию EURUSD, либо передайте как аргумент
DESTINATION="./data/${CURRENCY_PAIR}/"
MAX_RETRIES=3
CONCURRENT_JOBS=20

mkdir -p "$DESTINATION"

LOG_FILE="s3-sync-${CURRENCY_PAIR}-$(date +%Y%m%d_%H%M%S).log"

echo "Запуск синхронизации S3 для $CURRENCY_PAIR из $BUCKET_NAME в регионе $REGION..." | tee "$LOG_FILE"

for attempt in $(seq 1 $MAX_RETRIES); do
    echo "Попытка загрузки $attempt из $MAX_RETRIES..." | tee -a "$LOG_FILE"

    aws s3 sync "s3://$BUCKET_NAME/$CURRENCY_PAIR/" "$DESTINATION" \
      --region "$REGION" \
      --request-payer requester \
      --max-concurrent-requests 20 \
      --only-show-errors \
      --delete >> "$LOG_FILE" 2>&1

    EXIT_CODE=$?

    if [ $EXIT_CODE -eq 0 ]; then
        echo "✓ Загрузка успешно завершена $(date)" | tee -a "$LOG_FILE"
        echo "Всего файлов: $(find $DESTINATION -type f | wc -l)" | tee -a "$LOG_FILE"
        echo "Общий размер: $(du -sh $DESTINATION | cut -f1)" | tee -a "$LOG_FILE"
        exit 0
    else
        echo "✗ Загрузка не удалась, код выхода $EXIT_CODE. Повтор через 30 секунд..." | tee -a "$LOG_FILE"
        sleep 30
    fi
done

echo "✗ Загрузка не удалась после $MAX_RETRIES попыток" | tee -a "$LOG_FILE"
exit 1

Использование:

chmod +x download-price-history.sh

./download-price-history.sh EURUSD
./download-price-history.sh GBPUSD
./download-price-history.sh

7. Пакетная загрузка нескольких валютных пар

#!/bin/bash

BUCKET_NAME="cfg-public-proper-wallaby"
REGION="eu-west-1"
MAX_RETRIES=3

PAIRS=("EURUSD" "GBPUSD" "USDJPY" "AUDUSD" "NZDUSD" "USDCAD")

for PAIR in "${PAIRS[@]}"; do
    echo "Запуск загрузки для $PAIR..."

    DESTINATION="./data/${PAIR}/"
    mkdir -p "$DESTINATION"

    for attempt in $(seq 1 $MAX_RETRIES); do
        aws s3 sync "s3://$BUCKET_NAME/$PAIR/" "$DESTINATION" \
          --region "$REGION" \
          --request-payer requester \
          --max-concurrent-requests 20 \
          --only-show-errors \
          --delete

        if [ $? -eq 0 ]; then
            echo "✓ $PAIR завершено"
            break
        else
            echo "⚠ $PAIR попытка $attempt не удалась, повтор..."
            sleep 30
        fi
    done
done

echo "✓ Все загрузки завершены"

8. Использование Python Boto3 (программный доступ)

import boto3
import concurrent.futures
from botocore.config import Config
import logging
from pathlib import Path

logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger(__name__)

config = Config(
    max_pool_connections=20,
    retries={'max_attempts': 3, 'mode': 'adaptive'},
    connect_timeout=5,
    read_timeout=60,
    region_name='eu-west-1'
)

s3_client = boto3.client('s3', config=config)
s3_resource = boto3.resource('s3', config=config)

def download_file(bucket, key, filename):
    """Загрузить файл с включённым Requester Pays"""
    try:
        Path(filename).parent.mkdir(parents=True, exist_ok=True)
        s3_client.download_file(
            Bucket=bucket,
            Key=key,
            Filename=filename,
            ExtraArgs={'RequestPayer': 'requester'}
        )
        logger.info(f"✓ Загружено: {key}")
        return True
    except Exception as e:
        logger.error(f"✗ Ошибка загрузки {key}: {e}")
        return False

def batch_download(bucket, currency_pair, destination, max_workers=10):
    """Загрузить все файлы для валютной пары параллельно"""
    try:
        bucket_obj = s3_resource.Bucket(bucket)

        # Слэш в конце префикса предотвращает случайное совпадение с другим
        # инструментом, который может иметь этот же текст в качестве префикса
        # (например, "EURUSD" и гипотетический "EURUSDT").
        prefix = f"{currency_pair}/"
        objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
        logger.info(f"Найдено {len(objects)} объектов для загрузки для {currency_pair}")

        with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
            futures = []

            for obj in objects:
                if obj.key.endswith('/'):  # Пропустить директории
                    continue

                filename = f"{destination}/{obj.key}"
                futures.append(executor.submit(download_file, bucket, obj.key, filename))

            completed = 0
            failed = 0
            for future in concurrent.futures.as_completed(futures):
                if future.result():
                    completed += 1
                else:
                    failed += 1

            logger.info(f"\n=== Итоги загрузки ===")
            logger.info(f"Валютная пара: {currency_pair}")
            logger.info(f"Всего файлов: {len(futures)}")
            logger.info(f"Успешно: {completed}")
            logger.info(f"Не удалось: {failed}")

            return completed, failed

    except Exception as e:
        logger.error(f"Ошибка пакетной загрузки: {e}")
        return 0, len(objects)

# Использование - Одна валютная пара
if __name__ == "__main__":
    BUCKET = 'cfg-public-proper-wallaby'
    CURRENCY_PAIR = 'EURUSD'  # Измените при необходимости
    DESTINATION = f'./data/{CURRENCY_PAIR}'

    logger.info(f"Запуск загрузки из s3://{BUCKET}/{CURRENCY_PAIR}/")
    logger.info(f"Регион: eu-west-1")
    logger.info(f"Назначение: {DESTINATION}")

    completed, failed = batch_download(
        bucket=BUCKET,
        currency_pair=CURRENCY_PAIR,
        destination=DESTINATION,
        max_workers=10
    )

    if failed == 0:
        logger.info("\n✓ Все файлы успешно загружены!")
    else:
        logger.warning(f"\n⚠ {failed} файлов не удалось загрузить")

Использование:

python3 download_price_history.py
# Отредактируйте переменную CURRENCY_PAIR для других пар

Примечание: Учитывая, что для EUR/USD подтверждено 26 586 файлов, вызов list(bucket_obj.objects.filter(...)) выше перечислит их все через постраничные вызовы ListObjectsV2 перед началом загрузок. Сам этот шаг листинга влечёт собственные (небольшие) расходы на запросы, тарифицируемые отдельно от GET-запросов. Для очень больших префиксов рассмотрите ручную постраничную обработку или тестирование с частичной выборкой.

9. Загрузка всех валютных пар (Python)

import boto3
import concurrent.futures
from botocore.config import Config
import logging
from pathlib import Path

logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger(__name__)

config = Config(
    max_pool_connections=20,
    retries={'max_attempts': 3, 'mode': 'adaptive'},
    connect_timeout=5,
    read_timeout=60,
    region_name='eu-west-1'
)

s3_client = boto3.client('s3', config=config)
s3_resource = boto3.resource('s3', config=config)

def download_file(bucket, key, filename):
    """Загрузить файл с включённым Requester Pays"""
    try:
        Path(filename).parent.mkdir(parents=True, exist_ok=True)
        s3_client.download_file(
            Bucket=bucket,
            Key=key,
            Filename=filename,
            ExtraArgs={'RequestPayer': 'requester'}
        )
        return True
    except Exception as e:
        logger.error(f"✗ Ошибка загрузки {key}: {e}")
        return False

def download_all_pairs(bucket, destination_base, max_workers=10):
    """Загрузить всю историю для всех валютных пар"""
    try:
        bucket_obj = s3_resource.Bucket(bucket)

        pairs = set()
        for obj in bucket_obj.objects.filter(RequestPayer='requester'):
            pair = obj.key.split('/')[0]
            if pair and not obj.key.endswith('/'):
                pairs.add(pair)

        logger.info(f"Найдено {len(pairs)} валютных пар: {sorted(pairs)}")

        total_completed = 0
        total_failed = 0

        for pair in sorted(pairs):
            logger.info(f"\n=== Обработка {pair} ===")

            # Слэш в конце ограничивает выборку ключами именно этого инструмента
            prefix = f"{pair}/"
            pair_objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
            logger.info(f"Найдено {len(pair_objects)} объектов для {pair}")

            with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
                futures = []

                for obj in pair_objects:
                    if obj.key.endswith('/'):
                        continue

                    filename = f"{destination_base}/{obj.key}"
                    futures.append(executor.submit(download_file, bucket, obj.key, filename))

                completed = sum(1 for f in concurrent.futures.as_completed(futures) if f.result())
                failed = len(futures) - completed

                total_completed += completed
                total_failed += failed

                logger.info(f"✓ {pair}: {completed}/{len(futures)} файлов загружено")

        logger.info(f"\n=== ИТОГОВАЯ СВОДКА ===")
        logger.info(f"Всего завершено: {total_completed}")
        logger.info(f"Всего не удалось: {total_failed}")

        return total_completed, total_failed

    except Exception as e:
        logger.error(f"Ошибка: {e}")
        return 0, 0

# Использование
if __name__ == "__main__":
    BUCKET = 'cfg-public-proper-wallaby'
    DESTINATION_BASE = './data'

    logger.info(f"Запуск загрузки полного архива из s3://{BUCKET}/")
    logger.info(f"Регион: eu-west-1")
    logger.info(f"Назначение: {DESTINATION_BASE}")

    download_all_pairs(
        bucket=BUCKET,
        destination_base=DESTINATION_BASE,
        max_workers=10
    )

10. Краткий справочник команд

# Список всех доступных валютных пар
aws s3 ls s3://cfg-public-proper-wallaby/ \
  --region eu-west-1 \
  --request-payer requester

# Загрузка истории EURUSD
aws s3 sync s3://cfg-public-proper-wallaby/EURUSD/ ./EURUSD/ \
  --region eu-west-1 \
  --request-payer requester \
  --max-concurrent-requests 25

# Загрузка истории GBPUSD
aws s3 sync s3://cfg-public-proper-wallaby/GBPUSD/ ./GBPUSD/ \
  --region eu-west-1 \
  --request-payer requester \
  --max-concurrent-requests 25

# Подсчёт файлов в EURUSD
aws s3 ls s3://cfg-public-proper-wallaby/EURUSD/ \
  --region eu-west-1 \
  --request-payer requester \
  --recursive | wc -l

# Общий размер EURUSD
aws s3 ls s3://cfg-public-proper-wallaby/EURUSD/ \
  --region eu-west-1 \
  --request-payer requester \
  --recursive --summarize | grep "Total Size"

Ключевые обоснования

Повышение скорости ⚡

  • Параллельная обработка: 20-25 одновременных соединений максимизируют пропускную способность
  • Оптимизация региона EU-West-1: отсутствие межрегиональной задержки
  • Адаптивная стратегия повторов: интеллектуальное восстановление после сбоев
  • Пул соединений: минимизирует накладные расходы между запросами
  • Контроль пропускной способности: предотвращает перегрузку сети и проблемы со стабильностью

Повышение стабильности

  • Механизм повторных попыток: 3 попытки с интервалом 30 секунд
  • Логирование ошибок: подробные журналы для устранения неполадок
  • Таймаут соединения: 5 секунд на подключение, 60 секунд на чтение
  • Адаптивный режим: корректирует стратегию повторов в зависимости от типа ошибки
  • Проверки целостности: проверяют загруженные файлы
  • Плавная деградация: продолжает работу при частичных сбоях

Оценка стоимости

Структура цен

Цены AWS S3 на GET-запросы (Requester Pays, eu-west-1):

  • $0.0004 за 1000 запросов
  • $0.02 за ГБ передачи данных (тариф исходящего трафика eu-west-1)

Этапы декодирования и работы конвейера из разделов «Декодирование файлов .bi5» и «Полный конвейер» выполняются целиком на вашем локальном компьютере — они не добавляют дополнительных расходов AWS сверх стоимости загрузки, указанной здесь.

Оценка полного архива

⚠️ Полный архив: ~400 ГБ, ~20 000 000 файлов

Стоимость запросов  = (20 000 000 / 1 000) × $0.0004 = $8.00
Стоимость передачи = 400 × $0.02                      = $8.00
─────────────────────────────────────────────────────────
ИТОГО                                                  = $16.00

Средний размер файла по всему архиву: 400 ГБ / 20 000 000 файлов ≈ 20.97 КБ/файл

Реальный пример: EUR/USD (проверенные фактические данные)

Подтверждено через aws s3 ls --summarize:

Всего объектов: 26 586
Общий размер:    2.6 ГБ
Средний размер файла: 2.6 ГБ / 26 586 ≈ 100.1 КБ/файл

Стоимость запросов:

Количество запросов: 26 586
Стоимость за 1000 запросов: $0.0004

Общая стоимость запросов = (26 586 / 1 000) × $0.0004
                   = 26.586 × $0.0004
                   ≈ $0.0106

Стоимость передачи данных:

Размер данных: 2.6 ГБ
Стоимость за ГБ: $0.02
Общая стоимость передачи = 2.6 × $0.02 = $0.052

Общая стоимость загрузки EUR/USD (проверено):

Стоимость запросов:        $0.0106
Стоимость передачи данных:  $0.052
────────────────────────────
ИТОГО:                ~$0.06

Таблица сравнения стоимости

Сценарий Файлы Размер Стоимость запросов Стоимость передачи Итоговая стоимость
Полный архив 20М 400 ГБ $8.00 $8.00 $16.00
EUR/USD (проверено фактически) 26 586 2.6 ГБ $0.0106 $0.052 ~$0.06
5 пар (если аналогично EURUSD) ~133 000 ~13 ГБ $0.053 $0.26 ~$0.31
10 пар (если аналогично EURUSD) ~266 000 ~26 ГБ $0.106 $0.52 ~$0.63

Оценки для «5 пар» и «10 пар» предполагают количество файлов/размер, аналогичные EUR/USD. Фактические затраты будут варьироваться по парам — у основных пар может быть больше истории/файлов, чем у EUR/USD, у экзотических, вероятно, меньше.

Советы по оптимизации затрат

  • ✓ Загружайте только необходимые валютные пары (а не весь архив)
  • Всегда проверяйте фактический размер/количество файлов по каждой паре с помощью --summarize — реальный средний размер файла EUR/USD (~100 КБ) примерно в 5 раз больше среднего по всему архиву (~21 КБ), что показывает, насколько усреднённые значения могут вводить в заблуждение
  • ✓ Группируйте загрузки, чтобы минимизировать накладные расходы на запросы
  • ✓ Используйте флаг --delete, чтобы избежать повторной загрузки уже существующих файлов
  • ✓ При ~$0.06 за пару (проверено на EUR/USD) загрузка по отдельности даже десятков пар остаётся лишь небольшой долей от стоимости полного архива в $16
  • ✓ Используйте раздел 3 (Определение доступных инструментов), чтобы подтвердить точные названия пар перед запуском пакетных скриптов

Рекомендуемое действие перед массовыми загрузками

aws s3 ls s3://cfg-public-proper-wallaby/<PAIR>/ \
  --region eu-west-1 \
  --request-payer requester \
  --recursive --summarize | tail -3

Это вернёт точные значения Total Objects и Total Size, позволяя рассчитать точную стоимость по формуле:

Стоимость запросов  = (файлы / 1000) × $0.0004
Стоимость передачи = размер_в_ГБ × $0.02

Декодирование файлов .bi5 в читаемые тиковые данные

Файлы .bi5 Dukascopy представляют собой бинарные тиковые файлы, сжатые LZMA. При текущей структуре бакета каждый файл содержит один полный день тиковых данных для инструмента.

Структура файлов

Сжатие: необработанный поток LZMA (не формат-контейнер .xz)

Соглашение о путях (по дням):

SYMBOL/YEAR/MONTH/DAY_ticks.bi5

Пример: EURUSD/2024/00/15_ticks.bi5 → тики EUR/USD за 15 января 2024

⚠️ Месяц нумеруется с нуля (январь = 00, декабрь = 11).

Пустых файлов нет: отсутствие файла за конкретный день означает, что в этот день тики не фиксировались (выходные, праздники). Обрабатывайте отсутствующие ключи S3 / FileNotFoundError как «нет данных», а не как ошибку.

Формат распакованной записи

Каждый тик представляет собой фиксированную бинарную запись из 20 байт в формате big-endian:

Байты Поле Тип Примечания
0–3 Временная метка uint32 Миллисекунды от начала дня (UTC)
4–7 Цена Ask uint32 Требует масштабирования по точечному значению
8–11 Цена Bid uint32 Требует масштабирования по точечному значению
12–15 Объём Ask float32 В миллионах единиц базовой валюты
16–19 Объём Bid float32 В миллионах единиц базовой валюты

Масштабирование цены («точечное значение»)

Тип инструмента Точечное значение Пример
Большинство пар Forex 100 000 EUR/USD, GBP/USD
Пары с JPY 1 000 USD/JPY, EUR/JPY
Индексы/сырьевые товары Различается Проверяйте для каждого инструмента

⚠️ Вспомогательные функции ниже различают только «пара с JPY» и «всё остальное (100 000)». Для индексов, сырьевых товаров или любого не-Forex инструмента уточните правильное точечное значение у поставщика данных перед декодированием — молчаливое применение 100 000 к не-Forex инструменту приведёт к неверным ценам без какой-либо ошибки. Обновлённая функция get_point_value() в этом руководстве теперь выводит предупреждение, когда происходит откат к значению по умолчанию для нераспознанного, не-JPY инструмента, чтобы это не проходило незамеченным.

Python-декодер (один день)

import lzma
import struct
from datetime import datetime, timedelta

def decode_bi5_daily(filepath, day_start, point_value=100000):
    """
    Декодирует дневной тиковый файл .bi5.

    :param filepath: Путь к файлу .bi5
    :param day_start: datetime, представляющий 00:00 UTC этого дня
    :param point_value: 100000 для большинства пар, 1000 для пар с JPY
    :return: Список словарей с тиками
    """
    with open(filepath, 'rb') as f:
        compressed = f.read()

    raw = lzma.decompress(compressed)

    ticks = []
    record_size = 20
    for i in range(0, len(raw), record_size):
        chunk = raw[i:i+record_size]
        ms, ask, bid, ask_vol, bid_vol = struct.unpack('>IIIff', chunk)

        timestamp = day_start + timedelta(milliseconds=ms)
        ticks.append({
            'time': timestamp,
            'ask': ask / point_value,
            'bid': bid / point_value,
            'ask_volume': ask_vol,
            'bid_volume': bid_vol
        })

    return ticks

# Пример использования
day_start = datetime(2024, 1, 15, 0, 0, 0)
ticks = decode_bi5_daily('EURUSD/2024/00/15_ticks.bi5', day_start, point_value=100000)

print(f"Всего тиков за день: {len(ticks)}")
for tick in ticks[:5]:
    print(tick)

Скрипт пакетного декодирования (полная папка инструмента → CSV)

import lzma
import struct
import csv
import re
import logging
from datetime import datetime, timedelta
from pathlib import Path

logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')
logger = logging.getLogger(__name__)

JPY_PAIRS = {'USDJPY', 'EURJPY', 'GBPJPY', 'AUDJPY', 'NZDJPY', 'CADJPY', 'CHFJPY'}

# Известные не-Forex инструменты (индексы/сырьевые товары) с подтверждённым точечным значением.
# Дополняйте этот список по мере подтверждения других инструментов у поставщика данных.
OTHER_POINT_VALUES = {
    # 'US30': 100,
    # 'XAUUSD': 100,
}

def get_point_value(instrument):
    instrument = instrument.upper()
    if instrument in JPY_PAIRS:
        return 1000
    if instrument in OTHER_POINT_VALUES:
        return OTHER_POINT_VALUES[instrument]
    if len(instrument) != 6 or not instrument.isalpha():
        logger.warning(
            f"'{instrument}' не похож на стандартную 6-буквенную пару Forex и не имеет "
            f"подтверждённого точечного значения — используется значение по умолчанию 100000. Проверьте это, прежде чем доверять результату."
        )
    return 100000

def decode_bi5_daily(filepath, day_start, point_value):
    with open(filepath, 'rb') as f:
        compressed = f.read()

    raw = lzma.decompress(compressed)

    ticks = []
    record_size = 20
    for i in range(0, len(raw), record_size):
        chunk = raw[i:i+record_size]
        ms, ask, bid, ask_vol, bid_vol = struct.unpack('>IIIff', chunk)
        timestamp = day_start + timedelta(milliseconds=ms)
        ticks.append((timestamp, ask / point_value, bid / point_value, ask_vol, bid_vol))

    return ticks

def batch_decode_instrument(instrument_dir, output_csv):
    """
    Декодирует все дневные файлы .bi5 для инструмента в единый CSV-файл.
    Ожидает структуру папок: SYMBOL/YEAR/MONTH/DAY_ticks.bi5
    """
    instrument_dir = Path(instrument_dir)
    instrument_name = instrument_dir.name.upper()
    point_value = get_point_value(instrument_name)

    pattern = re.compile(r'(\d{2})_ticks\.bi5$')

    all_ticks = []
    file_count = 0
    error_count = 0

    for year_dir in sorted(instrument_dir.iterdir()):
        if not year_dir.is_dir():
            continue
        for month_dir in sorted(year_dir.iterdir()):
            if not month_dir.is_dir():
                continue

            year = int(year_dir.name)
            month = int(month_dir.name) + 1

            for bi5_file in sorted(month_dir.glob('*_ticks.bi5')):
                match = pattern.search(bi5_file.name)
                if not match:
                    continue

                day = int(match.group(1))
                day_start = datetime(year, month, day)

                try:
                    ticks = decode_bi5_daily(bi5_file, day_start, point_value)
                    all_ticks.extend(ticks)
                    file_count += 1
                except Exception as e:
                    print(f"✗ Ошибка декодирования {bi5_file}: {e}")
                    error_count += 1

    all_ticks.sort(key=lambda t: t[0])

    with open(output_csv, 'w', newline='') as f:
        writer = csv.writer(f)
        writer.writerow(['timestamp', 'ask', 'bid', 'ask_volume', 'bid_volume'])
        for tick in all_ticks:
            writer.writerow([tick[0].isoformat(), tick[1], tick[2], tick[3], tick[4]])

    print(f"\n=== Сводка декодирования ===")
    print(f"Инструмент: {instrument_name}")
    print(f"Обработано файлов: {file_count}")
    print(f"Ошибок: {error_count}")
    print(f"Всего тиков: {len(all_ticks)}")
    print(f"Результат: {output_csv}")

if __name__ == "__main__":
    batch_decode_instrument('./data/EURUSD', './EURUSD_ticks.csv')

Использование:

python3 decode_bi5_batch.py

Распространённые ошибки

  • ⚠️ Месяц нумеруется с нуля в пути S3, но не в вычислениях с датами — преобразуйте это явно.
  • ⚠️ Неверное точечное значение: несоответствие 100 000 и 1 000 незаметно искажает цены. Не-Forex инструменты (индексы/сырьевые товары) вообще не охватываются этим правилом — уточняйте их точечное значение явно (см. OTHER_POINT_VALUES выше).
  • ⚠️ Отсутствующие файлы ≠ ошибки: считайте это «в этот день тиков не было», а не сбоем декодирования.
  • ⚠️ Необработанный поток LZMA: некоторым библиотекам требуется явный режим «raw» — стандартные распаковщики .xz дадут сбой.
  • ⚠️ Устаревшие почасовые файлы: более старые файлы могут использовать миллисекунды от начала часа вместо миллисекунд от начала дня. Ни один из декодеров в этом руководстве не определяет это автоматически — оба предполагают текущий дневной формат (SYMBOL/YEAR/MONTH/DAY_ticks.bi5) повсеместно. Если вы работаете с инструментом или диапазоном дат, достаточно старым, чтобы предшествовать дневной структуре, проверьте фактическую временную базу файла перед декодированием, вместо того чтобы предполагать её соответствие описанному здесь дневному формату.

Полный конвейер (загрузка → декодирование → экспорт)

Объединяет загрузку, декодирование и экспорт в один автоматизированный скрипт — с опциональной фильтрацией по датам и очисткой.

Зависимости

pip install boto3 pandas pyarrow

pyarrow требуется только для экспорта в Parquet (рекомендуется для больших наборов данных).

Полный скрипт конвейера

#!/usr/bin/env python3
"""
Полный конвейер: загрузка -> декодирование -> экспорт
Дневные тиковые файлы .bi5 формата Dukascopy из S3 (Requester Pays)
"""

import argparse
import logging
import lzma
import struct
import shutil
from pathlib import Path
from datetime import datetime, timedelta

import boto3
import pandas as pd
from botocore.config import Config

BUCKET = 'cfg-public-proper-wallaby'
REGION = 'eu-west-1'

JPY_PAIRS = {'USDJPY', 'EURJPY', 'GBPJPY', 'AUDJPY', 'NZDJPY', 'CADJPY', 'CHFJPY'}

logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(levelname)s - %(message)s'
)
logger = logging.getLogger(__name__)

s3_config = Config(
    max_pool_connections=20,
    retries={'max_attempts': 3, 'mode': 'adaptive'},
    connect_timeout=5,
    read_timeout=60,
    region_name=REGION
)
s3_client = boto3.client('s3', config=s3_config)
s3_resource = boto3.resource('s3', config=s3_config)

def get_point_value(instrument):
    instrument = instrument.upper()
    if instrument in JPY_PAIRS:
        return 1000
    if len(instrument) != 6 or not instrument.isalpha():
        logger.warning(
            f"'{instrument}' не похож на стандартную 6-буквенную пару Forex — "
            f"используется значение по умолчанию 100000. Проверьте это, прежде чем доверять результату."
        )
    return 100000

# Шаг 1: Загрузка

def download_instrument(instrument, destination, start_date=None, end_date=None):
    """Загружает все файлы .bi5 для инструмента, опционально с фильтрацией по диапазону дат."""
    destination = Path(destination)
    destination.mkdir(parents=True, exist_ok=True)

    bucket_obj = s3_resource.Bucket(BUCKET)
    prefix = f"{instrument}/"

    downloaded = 0
    skipped = 0
    failed = 0

    logger.info(f"Получение списка объектов для {instrument}...")
    objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
    logger.info(f"Найдено {len(objects)} объектов")

    for obj in objects:
        if obj.key.endswith('/'):
            continue

        parts = obj.key.split('/')
        if len(parts) < 4:
            continue

        try:
            year = int(parts[1])
            month = int(parts[2]) + 1
            day = int(parts[3].split('_')[0])
            file_date = datetime(year, month, day)
        except (ValueError, IndexError):
            logger.warning(f"Не удалось разобрать дату из ключа: {obj.key}")
            continue

        if start_date and file_date < start_date:
            skipped += 1
            continue
        if end_date and file_date > end_date:
            skipped += 1
            continue

        local_path = destination / parts[1] / parts[2] / parts[3]
        local_path.parent.mkdir(parents=True, exist_ok=True)

        try:
            s3_client.download_file(
                Bucket=BUCKET,
                Key=obj.key,
                Filename=str(local_path),
                ExtraArgs={'RequestPayer': 'requester'}
            )
            downloaded += 1
        except Exception as e:
            logger.error(f"✗ Не удалось загрузить {obj.key}: {e}")
            failed += 1

    logger.info(f"Загрузка завершена: {downloaded} загружено, {skipped} пропущено (фильтр дат), {failed} не удалось")
    return downloaded, failed

# Шаг 2: Декодирование

def decode_bi5_daily(filepath, day_start, point_value):
    with open(filepath, 'rb') as f:
        compressed = f.read()

    raw = lzma.decompress(compressed)

    ticks = []
    record_size = 20
    for i in range(0, len(raw), record_size):
        chunk = raw[i:i + record_size]
        ms, ask, bid, ask_vol, bid_vol = struct.unpack('>IIIff', chunk)
        timestamp = day_start + timedelta(milliseconds=ms)
        ticks.append((timestamp, ask / point_value, bid / point_value, ask_vol, bid_vol))

    return ticks

def decode_instrument_folder(instrument_dir, instrument_name):
    """Декодирует все дневные файлы .bi5 в папке в список кортежей с тиками."""
    instrument_dir = Path(instrument_dir)
    point_value = get_point_value(instrument_name)

    all_ticks = []
    file_count = 0
    error_count = 0

    for year_dir in sorted(instrument_dir.iterdir()):
        if not year_dir.is_dir():
            continue
        for month_dir in sorted(year_dir.iterdir()):
            if not month_dir.is_dir():
                continue

            year = int(year_dir.name)
            month = int(month_dir.name) + 1

            for bi5_file in sorted(month_dir.glob('*_ticks.bi5')):
                try:
                    day = int(bi5_file.name.split('_')[0])
                    day_start = datetime(year, month, day)
                    ticks = decode_bi5_daily(bi5_file, day_start, point_value)
                    all_ticks.extend(ticks)
                    file_count += 1
                except Exception as e:
                    logger.error(f"✗ Ошибка декодирования {bi5_file}: {e}")
                    error_count += 1

    all_ticks.sort(key=lambda t: t[0])
    logger.info(f"Декодировано {file_count} файлов ({error_count} ошибок), всего тиков: {len(all_ticks)}")

    return all_ticks

# Шаг 3: Экспорт

def export_ticks(ticks, output_path, output_format='csv'):
    """Экспортирует декодированные тики в CSV или Parquet с использованием pandas."""
    df = pd.DataFrame(ticks, columns=['timestamp', 'ask', 'bid', 'ask_volume', 'bid_volume'])

    output_path = Path(output_path)
    output_path.parent.mkdir(parents=True, exist_ok=True)

    if output_format == 'csv':
        df.to_csv(output_path, index=False)
    elif output_format == 'parquet':
        df.to_parquet(output_path, index=False, engine='pyarrow', compression='snappy')
    else:
        raise ValueError(f"Unsupported format: {output_format}")

    logger.info(f"✓ Экспортировано {len(df)} тиков в {output_path} ({output_format})")
    return df

# Оркестрация конвейера

def run_pipeline(instrument, start_date=None, end_date=None,
                  output_format='parquet', cleanup=False,
                  raw_dir=None, output_dir='./output'):

    instrument = instrument.upper()
    raw_dir = raw_dir or f"./raw/{instrument}"
    output_path = Path(output_dir) / f"{instrument}_ticks.{output_format}"

    logger.info(f"=== Конвейер запущен для {instrument} ===")
    if start_date or end_date:
        logger.info(f"Диапазон дат: {start_date or 'самая ранняя'} по {end_date or 'самая поздняя'}")

    downloaded, failed = download_instrument(instrument, raw_dir, start_date, end_date)
    if failed:
        logger.warning(f"{failed} файл(ов) не удалось загрузить для {instrument} — продолжаем с {downloaded} успешно загруженными")
    if downloaded == 0:
        logger.warning(f"Для {instrument} не загружено ни одного файла — пропуск декодирования/экспорта")
        return

    ticks = decode_instrument_folder(raw_dir, instrument)
    if not ticks:
        logger.warning(f"Для {instrument} не декодировано ни одного тика")
        return

    export_ticks(ticks, output_path, output_format)

    if cleanup:
        logger.info(f"Очистка исходных файлов .bi5 в {raw_dir}...")
        shutil.rmtree(raw_dir, ignore_errors=True)
        logger.info("✓ Очистка завершена")

    logger.info(f"=== Конвейер завершён для {instrument} ===\n")

def parse_date(date_str):
    return datetime.strptime(date_str, '%Y-%m-%d')

if __name__ == "__main__":
    parser = argparse.ArgumentParser(description="Download, decode, and export Dukascopy .bi5 tick data")
    parser.add_argument('instrument', help="Currency pair, e.g. EURUSD")
    parser.add_argument('--start-date', type=parse_date, help="YYYY-MM-DD")
    parser.add_argument('--end-date', type=parse_date, help="YYYY-MM-DD")
    parser.add_argument('--format', choices=['csv', 'parquet'], default='parquet')
    parser.add_argument('--cleanup', action='store_true', help="Delete raw .bi5 files after export")
    parser.add_argument('--output-dir', default='./output')
    parser.add_argument('--raw-dir', default=None, help="Where to store raw .bi5 files (default: ./raw/<INSTRUMENT>)")

    args = parser.parse_args()

    run_pipeline(
        instrument=args.instrument,
        start_date=args.start_date,
        end_date=args.end_date,
        output_format=args.format,
        cleanup=args.cleanup,
        raw_dir=args.raw_dir,
        output_dir=args.output_dir
    )

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

# Загрузить, декодировать и экспортировать EURUSD в Parquet (по умолчанию)
python3 pipeline.py EURUSD

# Экспорт в CSV вместо этого
python3 pipeline.py EURUSD --format csv

# Фильтр по конкретному диапазону дат
python3 pipeline.py EURUSD --start-date 2024-01-01 --end-date 2024-03-31

# Очистка исходных файлов .bi5 после экспорта
python3 pipeline.py EURUSD --cleanup

# Пользовательский каталог для результатов
python3 pipeline.py GBPUSD --output-dir ./exports --cleanup

Пакетный конвейер для нескольких инструментов

#!/bin/bash
# run_all_pipelines.sh

INSTRUMENTS=("EURUSD" "GBPUSD" "USDJPY" "AUDUSD")

for INSTRUMENT in "${INSTRUMENTS[@]}"; do
    echo "=== Обработка $INSTRUMENT ==="
    python3 pipeline.py "$INSTRUMENT" --format parquet --cleanup
done

echo "✓ Все конвейеры завершены"

Лучшие практики ✅

Загрузка и контроль затрат

  • ✓ Всегда указывайте флаги --region eu-west-1 и --request-payer requester
  • ✓ Проверяйте фактический размер/количество файлов по каждой паре с помощью --summarize перед массовыми загрузками — усреднённые значения могут вводить в заблуждение (например, реальный средний размер файла EUR/USD примерно в 5 раз больше среднего по всему архиву)
  • ✓ Используйте aws s3 sync (а не cp) для нескольких файлов — передаются только изменённые данные
  • ✓ Используйте флаг --delete, чтобы избежать повторной загрузки уже существующих файлов
  • ✓ Сначала запускайте определение инструментов, чтобы подтвердить точные, действительные названия пар — это позволяет избежать неудачных синхронизаций из-за опечаток
  • ✓ Ведите лог-файлы для аудита
  • ✓ Настройте CloudWatch-алармы для неудачных загрузок
  • ✓ Пересчитывайте затраты при существенном изменении размера бакета/количества файлов
  • ✓ По возможности отдавайте предпочтение загрузке конкретных пар, а не всего архива

Декодирование и целостность данных

  • ✓ Всегда подтверждайте правильное точечное значение (100 000 или 1 000) для каждого инструмента перед декодированием — неверное масштабирование незаметно искажает цены
  • ✓ Считайте отсутствующие дневные файлы .bi5 признаком «в этот день тиков не было» (выходные/праздники), а не ошибками
  • ✓ Никогда не смешивайте устаревшие почасовые файлы .bi5 с новыми дневными файлами в одном конвейере без предварительного определения формата
  • ✓ Проверяйте режим распаковки LZMA — некоторым библиотекам требуется явный режим «raw» (без заголовков контейнера)
  • ✓ Сортируйте декодированные тики в хронологическом порядке перед экспортом на случай несоответствий в порядке файловой системы

Конвейер и автоматизация

  • ✓ Используйте --cleanup, чтобы удалить исходные файлы .bi5 после декодирования — их всегда можно повторно загрузить из S3 при необходимости
  • ✓ Предпочитайте Parquet вместо CSV для тиковых данных — обычно в 5-10 раз меньше и быстрее при запросах
  • ✓ Используйте фильтры --start-date / --end-date при тестировании, чтобы избежать ненужных расходов на S3
  • ✓ Сочетайте с определением инструментов, чтобы автоматически перебирать все допустимые пары вместо жёстко заданных списков
  • ✓ Логируйте каждый запуск конвейера (количество загруженных/декодированных/экспортированных) для возможности аудита

Общее

  • ✓ Выполняйте загрузки в непиковые часы, чтобы снизить конкуренцию (не стоимость)
  • ✓ Внедряйте контрольные суммы для проверки целостности данных там, где это критично
  • ✓ Отслеживайте общий объём переданных данных со временем для контроля совокупных затрат