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 /(ouDelimiter='/'dans Boto3) est essentielle — elle déclenche une réponse allégéeCommonPrefixesau 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
--deletepour é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 /
FileNotFoundErrorcomme « 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_VALUESci-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
.xzstandards é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
pyarrown'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-1et--request-payer requester - ✓ Vérifier la taille/le nombre de fichiers réels par paire avec
--summarizeavant 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 noncp) pour plusieurs fichiers — seules les données modifiées sont transférées - ✓ Utiliser le drapeau
--deletepour é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
.bi5quotidiens 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
.bi5avec 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
--cleanuppour supprimer les fichiers.bi5bruts 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-datelors 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