Données historiques de prix

Aperçu

Ce guide vous permet d'accéder aux fichiers de données historiques de prix du bucket S3 cfg-public-proper-wallaby (eu-west-1), configuré avec la politique Requester Pays, garantissant des performances de téléchargement optimales et la stabilité du système.

Prérequis

  • AWS CLI installé et configuré
  • Identifiants AWS valides avec les autorisations appropriées
  • Accès à la région eu-west-1
  • Compréhension du modèle de facturation requester pays

Étapes de configuration

1. Définir les identifiants AWS et la région

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

2. Vérifier l'accès au bucket

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

3. Découvrir les instruments disponibles

Avant de télécharger, listez toutes les paires de devises (instruments) disponibles dans le bucket. Utilisez --delimiter / (ou simplement aws s3 ls) pour ne lister que les préfixes de premier niveau — évitez le listage récursif, qui parcourrait chaque fichier uniquement pour trouver les ~20-30 noms de dossiers.

Liste rapide (la plus simple) :

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

Liste propre (noms uniquement, triés, enregistrés dans un fichier) :

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 "Total des instruments : $(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):
    """Récupère tous les préfixes d'instruments de premier niveau (paires de devises)"""
    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"Total des instruments trouvés : {len(instruments)}\n")
    for inst in instruments:
        print(inst)

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

Utilisation :

python3 list_instruments.py

⚠️ Remarque sur les performances et les coûts : L'utilisation de --delimiter / (ou Delimiter='/' dans Boto3) est essentielle — elle déclenche une réponse allégée CommonPrefixes au lieu d'énumérer chaque fichier. Cela ne coûte que quelques requêtes LIST peu coûteuses et se termine en quelques secondes, contrairement à une analyse complète du bucket, lente et plus coûteuse, sans le délimiteur.

4. Commande d'accès S3 de base avec Requester Pays

# Télécharger un seul fichier
aws s3 cp s3://cfg-public-proper-wallaby/EURUSD/file.json . \
  --region eu-west-1 \
  --request-payer requester

5. Téléchargement par lots avec optimisation des performances

# Télécharger tout l'historique pour une paire de devises spécifique
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

Paramètres d'optimisation des performances

Paramètre Valeur Avantage
--max-concurrent-requests 20-30 Augmente les téléchargements parallèles
--max-bandwidth 100MB/s Évite les goulots d'étranglement réseau
--no-progress - Réduit la charge d'E/S
--region eu-west-1 Réduit la latence au sein de la région UE

6. Script de téléchargement avancé (vitesse et stabilité)

#!/bin/bash

BUCKET_NAME="cfg-public-proper-wallaby"
REGION="eu-west-1"
CURRENCY_PAIR="${1:-EURUSD}"  # EURUSD par défaut, ou transmis en argument
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 "Démarrage de la synchronisation S3 pour $CURRENCY_PAIR depuis $BUCKET_NAME dans la région $REGION..." | tee "$LOG_FILE"

for attempt in $(seq 1 $MAX_RETRIES); do
    echo "Tentative de téléchargement $attempt sur $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 "✓ Téléchargement terminé avec succès le $(date)" | tee -a "$LOG_FILE"
        echo "Total des fichiers : $(find $DESTINATION -type f | wc -l)" | tee -a "$LOG_FILE"
        echo "Taille totale : $(du -sh $DESTINATION | cut -f1)" | tee -a "$LOG_FILE"
        exit 0
    else
        echo "✗ Échec du téléchargement avec le code de sortie $EXIT_CODE. Nouvelle tentative dans 30 secondes..." | tee -a "$LOG_FILE"
        sleep 30
    fi
done

echo "✗ Échec du téléchargement après $MAX_RETRIES tentatives" | tee -a "$LOG_FILE"
exit 1

Utilisation :

chmod +x download-price-history.sh

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

7. Téléchargement par lots de plusieurs paires de devises

#!/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 "Démarrage du téléchargement pour $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 terminé"
            break
        else
            echo "⚠ $PAIR tentative $attempt échouée, nouvelle tentative..."
            sleep 30
        fi
    done
done

echo "✓ Tous les téléchargements sont terminés"

8. Utilisation de Python Boto3 (accès programmatique)

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):
    """Télécharge un fichier avec Requester Pays activé"""
    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"✓ Téléchargé : {key}")
        return True
    except Exception as e:
        logger.error(f"✗ Erreur lors du téléchargement de {key} : {e}")
        return False

def batch_download(bucket, currency_pair, destination, max_workers=10):
    """Télécharge tous les fichiers d'une paire de devises en parallèle"""
    try:
        bucket_obj = s3_resource.Bucket(bucket)

        # Une barre oblique finale sur le préfixe évite de correspondre par erreur
        # à un autre instrument qui aurait ce texte comme préfixe
        # (par ex. "EURUSD" face à un hypothétique "EURUSDT").
        prefix = f"{currency_pair}/"
        objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
        logger.info(f"{len(objects)} objets trouvés à télécharger pour {currency_pair}")

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

            for obj in objects:
                if obj.key.endswith('/'):  # Ignorer les répertoires
                    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=== Résumé du téléchargement ===")
            logger.info(f"Paire de devises : {currency_pair}")
            logger.info(f"Total des fichiers : {len(futures)}")
            logger.info(f"Réussis : {completed}")
            logger.info(f"Échoués : {failed}")

            return completed, failed

    except Exception as e:
        logger.error(f"Erreur de téléchargement par lots : {e}")
        return 0, len(objects)

# Utilisation - Paire de devises unique
if __name__ == "__main__":
    BUCKET = 'cfg-public-proper-wallaby'
    CURRENCY_PAIR = 'EURUSD'  # À modifier selon les besoins
    DESTINATION = f'./data/{CURRENCY_PAIR}'

    logger.info(f"Démarrage du téléchargement depuis s3://{BUCKET}/{CURRENCY_PAIR}/")
    logger.info(f"Région : eu-west-1")
    logger.info(f"Destination : {DESTINATION}")

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

    if failed == 0:
        logger.info("\n✓ Tous les fichiers ont été téléchargés avec succès !")
    else:
        logger.warning(f"\n⚠ {failed} fichiers n'ont pas pu être téléchargés")

Utilisation :

python3 download_price_history.py
# Modifier la variable CURRENCY_PAIR pour d'autres paires

Remarque : Étant donné que 26 586 fichiers ont été confirmés pour EUR/USD, l'appel list(bucket_obj.objects.filter(...)) ci-dessus les énumérera tous via des appels paginés ListObjectsV2 avant le début des téléchargements. Cette étape de listage entraîne elle-même son propre coût de requête (faible), facturé séparément des requêtes GET. Pour des préfixes très volumineux, envisagez une pagination manuelle ou un test avec des récupérations partielles.

9. Télécharger toutes les paires de devises (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):
    """Télécharge un fichier avec Requester Pays activé"""
    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"✗ Erreur lors du téléchargement de {key} : {e}")
        return False

def download_all_pairs(bucket, destination_base, max_workers=10):
    """Télécharge tout l'historique pour toutes les paires de devises"""
    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)} paires de devises trouvées : {sorted(pairs)}")

        total_completed = 0
        total_failed = 0

        for pair in sorted(pairs):
            logger.info(f"\n=== Traitement de {pair} ===")

            # La barre oblique finale limite précisément cette étape aux clés de cet instrument
            prefix = f"{pair}/"
            pair_objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
            logger.info(f"{len(pair_objects)} objets trouvés pour {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)} fichiers téléchargés")

        logger.info(f"\n=== RÉSUMÉ FINAL ===")
        logger.info(f"Total réussis : {total_completed}")
        logger.info(f"Total échoués : {total_failed}")

        return total_completed, total_failed

    except Exception as e:
        logger.error(f"Erreur : {e}")
        return 0, 0

# Utilisation
if __name__ == "__main__":
    BUCKET = 'cfg-public-proper-wallaby'
    DESTINATION_BASE = './data'

    logger.info(f"Démarrage du téléchargement de l'archive complète depuis s3://{BUCKET}/")
    logger.info(f"Région : eu-west-1")
    logger.info(f"Destination : {DESTINATION_BASE}")

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

10. Commandes de référence rapide

# Lister toutes les paires de devises disponibles
aws s3 ls s3://cfg-public-proper-wallaby/ \
  --region eu-west-1 \
  --request-payer requester

# Télécharger l'historique EURUSD
aws s3 sync s3://cfg-public-proper-wallaby/EURUSD/ ./EURUSD/ \
  --region eu-west-1 \
  --request-payer requester \
  --max-concurrent-requests 25

# Télécharger l'historique GBPUSD
aws s3 sync s3://cfg-public-proper-wallaby/GBPUSD/ ./GBPUSD/ \
  --region eu-west-1 \
  --request-payer requester \
  --max-concurrent-requests 25

# Compter les fichiers dans EURUSD
aws s3 ls s3://cfg-public-proper-wallaby/EURUSD/ \
  --region eu-west-1 \
  --request-payer requester \
  --recursive | wc -l

# Obtenir la taille totale d'EURUSD
aws s3 ls s3://cfg-public-proper-wallaby/EURUSD/ \
  --region eu-west-1 \
  --request-payer requester \
  --recursive --summarize | grep "Total Size"

Points clés de motivation

Amélioration de la vitesse ⚡

  • Traitement parallèle : 20-25 connexions simultanées maximisent le débit
  • Optimisation EU-West-1 : aucune latence inter-régionale
  • Stratégie de nouvelle tentative adaptative : récupération intelligente en cas d'échec
  • Regroupement de connexions : minimise la surcharge entre les requêtes
  • Contrôle de la bande passante : évite la saturation du réseau et les problèmes de stabilité

Amélioration de la stabilité

  • Mécanisme de nouvelle tentative : 3 tentatives avec des intervalles de 30 secondes
  • Journalisation des erreurs : journaux détaillés pour le dépannage
  • Délai d'expiration de connexion : 5 secondes pour la connexion, 60 secondes pour la lecture
  • Mode adaptatif : ajuste la stratégie de nouvelle tentative selon le type d'erreur
  • Contrôles d'intégrité : valident les téléchargements de fichiers
  • Dégradation gracieuse : continue en cas d'échecs partiels

Estimation des coûts

Structure tarifaire

Tarification des requêtes GET AWS S3 (Requester Pays, eu-west-1) :

  • 0,0004 $ pour 1 000 requêtes
  • 0,02 $ par Go de transfert de données (tarif de sortie eu-west-1)

Les étapes de décodage et de pipeline des sections « Décodage des fichiers .bi5 » et « Pipeline de bout en bout » s'exécutent entièrement sur votre machine locale — elles n'entraînent aucun coût AWS supplémentaire au-delà des coûts de téléchargement présentés ici.

Estimation de l'archive complète

⚠️ Archive complète : ~400 Go, ~20 000 000 fichiers

Coût des requêtes  = (20 000 000 / 1 000) × 0,0004 $ = 8,00 $
Coût du transfert = 400 × 0,02 $                      = 8,00 $
─────────────────────────────────────────────────────────
TOTAL                                                  = 16,00 $

Taille moyenne des fichiers à l'échelle de l'archive : 400 Go / 20 000 000 fichiers ≈ 20,97 Ko/fichier

Exemple concret : EUR/USD (données réelles vérifiées)

Confirmé via aws s3 ls --summarize :

Total des objets : 26 586
Taille totale :    2,6 Go
Taille moyenne des fichiers : 2,6 Go / 26 586 ≈ 100,1 Ko/fichier

Coûts des requêtes :

Nombre de requêtes : 26 586
Coût pour 1 000 requêtes : 0,0004 $

Coût total des requêtes = (26 586 / 1 000) × 0,0004 $
                   = 26,586 × 0,0004 $
                   ≈ 0,0106 $

Coûts de transfert de données :

Taille des données : 2,6 Go
Coût par Go : 0,02 $
Coût total du transfert = 2,6 × 0,02 $ = 0,052 $

Coût total du téléchargement EUR/USD (vérifié) :

Coûts des requêtes :        0,0106 $
Coûts de transfert de données : 0,052 $
────────────────────────────
TOTAL :                ~0,06 $

Tableau comparatif des coûts

Scénario Fichiers Taille Coût des requêtes Coût du transfert Coût total
Archive complète 20M 400 Go 8,00 $ 8,00 $ 16,00 $
EUR/USD (vérifié, réel) 26 586 2,6 Go 0,0106 $ 0,052 $ ~0,06 $
5 paires (si similaires à EURUSD) ~133 000 ~13 Go 0,053 $ 0,26 $ ~0,31 $
10 paires (si similaires à EURUSD) ~266 000 ~26 Go 0,106 $ 0,52 $ ~0,63 $

Les estimations pour « 5 paires » et « 10 paires » supposent un nombre de fichiers/une taille similaires à EUR/USD. Les coûts réels varieront selon la paire — les paires majeures pourraient avoir plus d'historique/de fichiers qu'EUR/USD, les paires exotiques probablement moins.

Conseils d'optimisation des coûts

  • ✓ Ne téléchargez que les paires de devises nécessaires (pas l'archive entière)
  • Vérifiez toujours la taille/le nombre de fichiers réels par paire avec --summarize — la taille moyenne réelle des fichiers d'EUR/USD (~100 Ko) est environ 5 fois supérieure à la moyenne à l'échelle de l'archive (~21 Ko), ce qui montre à quel point les moyennes peuvent être trompeuses
  • ✓ Regroupez les téléchargements pour minimiser la surcharge des requêtes
  • ✓ Utilisez le drapeau --delete pour éviter de retélécharger des fichiers déjà existants
  • ✓ À ~0,06 $ par paire (vérifié sur EUR/USD), télécharger individuellement des dizaines de paires ne représente qu'une infime fraction du coût de l'archive complète, soit 16 $
  • ✓ Utilisez la section 3 (Découvrir les instruments disponibles) pour confirmer les noms exacts des paires avant d'exécuter des scripts par lots

Action recommandée avant les téléchargements en masse

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

Cela renvoie les valeurs exactes de Total Objects et Total Size, permettant de calculer un coût précis via :

Coût des requêtes  = (fichiers / 1 000) × 0,0004 $
Coût du transfert = taille_en_Go × 0,02 $

Décodage des fichiers .bi5 en données tick lisibles

Les fichiers .bi5 de Dukascopy sont des fichiers binaires de ticks compressés en LZMA. Avec la structure actuelle du bucket, chaque fichier représente une journée complète de données tick par instrument.

Structure des fichiers

Compression : flux LZMA brut (pas au format conteneur .xz)

Convention de chemin (quotidienne) :

SYMBOL/YEAR/MONTH/DAY_ticks.bi5

Exemple : EURUSD/2024/00/15_ticks.bi5 → ticks EUR/USD du 15 janvier 2024

⚠️ Le mois est indexé à partir de zéro (janvier = 00, décembre = 11).

Aucun fichier vide : un fichier manquant pour un jour donné signifie qu'aucun tick n'a été enregistré ce jour-là (week-ends, jours fériés). Traitez les clés S3 manquantes / FileNotFoundError comme « aucune donnée », et non comme une erreur.

Format de l'enregistrement décompressé

Chaque tick est un enregistrement binaire fixe de 20 octets, en big-endian :

Octets Champ Type Remarques
0–3 Horodatage uint32 Millisecondes depuis le début de la journée (UTC)
4–7 Prix Ask uint32 Nécessite une mise à l'échelle par la valeur de point
8–11 Prix Bid uint32 Nécessite une mise à l'échelle par la valeur de point
12–15 Volume Ask float32 En millions d'unités de la devise de base
16–19 Volume Bid float32 En millions d'unités de la devise de base

Mise à l'échelle des prix (« valeur de point »)

Type d'instrument Valeur de point Exemple
La plupart des paires FX 100 000 EUR/USD, GBP/USD
Paires JPY 1 000 USD/JPY, EUR/JPY
Indices/matières premières Variable À vérifier par instrument

⚠️ Les fonctions d'aide ci-dessous ne distinguent que « paire JPY » de « tout le reste (100 000) ». Pour les indices, matières premières ou tout instrument non-FX, confirmez la valeur de point correcte auprès du fournisseur de données avant le décodage — appliquer silencieusement 100 000 à un instrument non-FX produira des prix erronés sans qu'aucune erreur ne soit levée. La fonction get_point_value() mise à jour dans ce guide émet désormais un avertissement lorsqu'elle revient à la valeur par défaut pour un instrument non reconnu et non-JPY, afin que cela ne passe plus inaperçu.

Décodeur Python (jour unique)

import lzma
import struct
from datetime import datetime, timedelta

def decode_bi5_daily(filepath, day_start, point_value=100000):
    """
    Décode un fichier tick .bi5 quotidien.

    :param filepath: Chemin vers le fichier .bi5
    :param day_start: datetime représentant 00:00 UTC de ce jour
    :param point_value: 100000 pour la plupart des paires, 1000 pour les paires JPY
    :return: Liste de dictionnaires de ticks
    """
    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

# Exemple d'utilisation
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"Total de ticks pour la journée : {len(ticks)}")
for tick in ticks[:5]:
    print(tick)

Script de décodage par lots (dossier d'instrument complet → 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'}

# Instruments non-FX connus (indices/matières premières) avec une valeur de point confirmée.
# Complétez cette liste au fur et à mesure que vous confirmez d'autres instruments auprès du fournisseur de données.
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}' ne ressemble pas à une paire FX standard à 6 lettres et n'a pas "
            f"de valeur de point confirmée — utilisation de la valeur par défaut 100000. Vérifiez ceci avant de faire confiance au résultat."
        )
    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):
    """
    Décode tous les fichiers .bi5 quotidiens d'un instrument en un seul CSV.
    Attend la structure de dossiers : 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"✗ Erreur lors du décodage de {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=== Résumé du décodage ===")
    print(f"Instrument : {instrument_name}")
    print(f"Fichiers traités : {file_count}")
    print(f"Erreurs : {error_count}")
    print(f"Total de ticks : {len(all_ticks)}")
    print(f"Sortie : {output_csv}")

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

Utilisation :

python3 decode_bi5_batch.py

Pièges courants

  • ⚠️ Le mois est indexé à partir de zéro dans le chemin S3, mais pas dans les calculs de date — à convertir explicitement.
  • ⚠️ Valeur de point erronée : une confusion entre 100 000 et 1 000 corrompt silencieusement les prix. Les instruments non-FX (indices/matières premières) ne sont absolument pas couverts par cette règle — confirmez explicitement leur valeur de point (voir OTHER_POINT_VALUES ci-dessus).
  • ⚠️ Fichiers manquants ≠ erreurs : à traiter comme « aucun tick ce jour-là », pas comme un échec de décodage.
  • ⚠️ Flux LZMA brut : certaines bibliothèques nécessitent un mode « raw » explicite — les décompresseurs .xz standards échoueront.
  • ⚠️ Anciens fichiers horaires : les fichiers plus anciens peuvent utiliser des millisecondes depuis le début de l'heure au lieu de millisecondes depuis le début de la journée. Aucun des décodeurs de ce guide ne détecte cela automatiquement — les deux supposent le format quotidien actuel (SYMBOL/YEAR/MONTH/DAY_ticks.bi5) de bout en bout. Si vous travaillez avec un instrument ou une plage de dates suffisamment ancienne pour précéder la disposition quotidienne, vérifiez la base temporelle réelle du fichier avant de décoder, plutôt que de supposer qu'elle correspond au format quotidien décrit ici.

Pipeline de bout en bout (téléchargement → décodage → export)

Combine le téléchargement, le décodage et l'export en un seul script automatisé — avec filtrage optionnel par date et nettoyage.

Dépendances

pip install boto3 pandas pyarrow

pyarrow n'est requis que pour l'export au format Parquet (recommandé pour les grands ensembles de données).

Script de pipeline complet

#!/usr/bin/env python3
"""
Pipeline de bout en bout : téléchargement -> décodage -> export
Fichiers tick quotidiens .bi5 de type Dukascopy depuis 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}' ne ressemble pas à une paire FX standard à 6 lettres — "
            f"utilisation de la valeur de point par défaut 100000. Vérifiez ceci avant de faire confiance au résultat."
        )
    return 100000

# Étape 1 : Téléchargement

def download_instrument(instrument, destination, start_date=None, end_date=None):
    """Télécharge tous les fichiers .bi5 d'un instrument, éventuellement filtrés par plage de dates."""
    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"Listage des objets pour {instrument}...")
    objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
    logger.info(f"{len(objects)} objets trouvés")

    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"Impossible d'extraire la date de la clé : {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"✗ Échec du téléchargement de {obj.key} : {e}")
            failed += 1

    logger.info(f"Téléchargement terminé : {downloaded} téléchargés, {skipped} ignorés (filtre de date), {failed} échoués")
    return downloaded, failed

# Étape 2 : Décodage

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):
    """Décode tous les fichiers .bi5 quotidiens d'un dossier en une liste de tuples de ticks."""
    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"✗ Erreur lors du décodage de {bi5_file} : {e}")
                    error_count += 1

    all_ticks.sort(key=lambda t: t[0])
    logger.info(f"{file_count} fichiers décodés ({error_count} erreurs), {len(all_ticks)} ticks au total")

    return all_ticks

# Étape 3 : Export

def export_ticks(ticks, output_path, output_format='csv'):
    """Exporte les ticks décodés en CSV ou Parquet à l'aide de 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)} ticks exportés vers {output_path} ({output_format})")
    return df

# Orchestration du pipeline

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"=== Pipeline démarré pour {instrument} ===")
    if start_date or end_date:
        logger.info(f"Plage de dates : {start_date or 'la plus ancienne'} à {end_date or 'la plus récente'}")

    downloaded, failed = download_instrument(instrument, raw_dir, start_date, end_date)
    if failed:
        logger.warning(f"{failed} fichier(s) n'ont pas pu être téléchargés pour {instrument} — poursuite avec les {downloaded} téléchargés avec succès")
    if downloaded == 0:
        logger.warning(f"Aucun fichier téléchargé pour {instrument} — décodage/export ignoré")
        return

    ticks = decode_instrument_folder(raw_dir, instrument)
    if not ticks:
        logger.warning(f"Aucun tick décodé pour {instrument}")
        return

    export_ticks(ticks, output_path, output_format)

    if cleanup:
        logger.info(f"Nettoyage des fichiers .bi5 bruts dans {raw_dir}...")
        shutil.rmtree(raw_dir, ignore_errors=True)
        logger.info("✓ Nettoyage terminé")

    logger.info(f"=== Pipeline terminé pour {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
    )

Exemples d'utilisation

# Télécharger, décoder et exporter EURUSD vers Parquet (par défaut)
python3 pipeline.py EURUSD

# Exporter en CSV à la place
python3 pipeline.py EURUSD --format csv

# Filtrer sur une plage de dates spécifique
python3 pipeline.py EURUSD --start-date 2024-01-01 --end-date 2024-03-31

# Nettoyer les fichiers .bi5 bruts après l'export
python3 pipeline.py EURUSD --cleanup

# Répertoire de sortie personnalisé
python3 pipeline.py GBPUSD --output-dir ./exports --cleanup

Pipeline par lots pour plusieurs instruments

#!/bin/bash
# run_all_pipelines.sh

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

for INSTRUMENT in "${INSTRUMENTS[@]}"; do
    echo "=== Traitement de $INSTRUMENT ==="
    python3 pipeline.py "$INSTRUMENT" --format parquet --cleanup
done

echo "✓ Tous les pipelines sont terminés"

Bonnes pratiques ✅

Téléchargement et contrôle des coûts

  • ✓ Toujours inclure les drapeaux --region eu-west-1 et --request-payer requester
  • ✓ Vérifier la taille/le nombre de fichiers réels par paire avec --summarize avant les téléchargements en masse — les moyennes peuvent être trompeuses (par ex. la taille moyenne réelle des fichiers EUR/USD est ~5 fois supérieure à la moyenne à l'échelle de l'archive)
  • ✓ Utiliser aws s3 sync (et non cp) pour plusieurs fichiers — seules les données modifiées sont transférées
  • ✓ Utiliser le drapeau --delete pour éviter de retélécharger des fichiers déjà existants
  • ✓ Exécuter d'abord la découverte des instruments pour confirmer des noms de paires exacts et valides — évite les échecs de synchronisation dus à des fautes de frappe
  • ✓ Conserver les fichiers journaux pour les pistes d'audit
  • ✓ Configurer des alarmes CloudWatch pour les téléchargements échoués
  • ✓ Recalculer les coûts chaque fois que la taille du bucket/le nombre de fichiers change de manière significative
  • ✓ Privilégier le téléchargement de paires spécifiques plutôt que l'archive entière lorsque c'est possible

Décodage et intégrité des données

  • ✓ Toujours confirmer la valeur de point correcte (100 000 contre 1 000) par instrument avant le décodage — une mauvaise mise à l'échelle corrompt silencieusement les prix
  • ✓ Traiter les fichiers .bi5 quotidiens manquants comme « aucun tick ce jour-là » (week-ends/jours fériés), pas comme des erreurs
  • ✓ Ne jamais mélanger d'anciens fichiers horaires .bi5 avec de nouveaux fichiers quotidiens dans le même pipeline sans détecter d'abord le format
  • ✓ Valider le mode de décompression LZMA — certaines bibliothèques nécessitent un mode « raw » explicite (sans en-têtes de conteneur)
  • ✓ Trier les ticks décodés par ordre chronologique avant l'export, en cas d'incohérences dans l'ordre du système de fichiers

Pipeline et automatisation

  • ✓ Utiliser --cleanup pour supprimer les fichiers .bi5 bruts après le décodage — ils peuvent toujours être retéléchargés depuis S3 si nécessaire
  • ✓ Privilégier Parquet plutôt que CSV pour les données au niveau tick — généralement 5 à 10 fois plus petit et plus rapide à interroger
  • ✓ Utiliser les filtres --start-date / --end-date lors des tests, pour éviter des coûts S3 inutiles
  • ✓ Combiner avec la découverte des instruments pour parcourir automatiquement toutes les paires valides au lieu de coder des listes en dur
  • ✓ Journaliser chaque exécution du pipeline (nombres téléchargés/décodés/exportés) à des fins d'audit

Général

  • ✓ Exécuter les téléchargements en dehors des heures de pointe pour réduire la contention (pas le coût)
  • ✓ Implémenter des sommes de contrôle pour la vérification de l'intégrité des données là où c'est critique
  • ✓ Surveiller le volume total de données transférées dans le temps pour contrôler les coûts cumulés

common.disclaimer