Исторические ценовые данные
Обзор
Это руководство поможет вам получить доступ к файлам исторических ценовых данных из 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 - ✓ Сочетайте с определением инструментов, чтобы автоматически перебирать все допустимые пары вместо жёстко заданных списков
- ✓ Логируйте каждый запуск конвейера (количество загруженных/декодированных/экспортированных) для возможности аудита
Общее
- ✓ Выполняйте загрузки в непиковые часы, чтобы снизить конкуренцию (не стоимость)
- ✓ Внедряйте контрольные суммы для проверки целостности данных там, где это критично
- ✓ Отслеживайте общий объём переданных данных со временем для контроля совокупных затрат