Historische Preisdaten
Überblick
Diese Anleitung ermöglicht Ihnen den Zugriff auf Dateien mit historischen Preisdaten aus dem S3-Bucket cfg-public-proper-wallaby (eu-west-1), der mit der Requester-Pays-Richtlinie konfiguriert ist, um eine optimale Download-Leistung und Systemstabilität zu gewährleisten.
Voraussetzungen
- Installiertes und konfiguriertes AWS CLI
- Gültige AWS-Zugangsdaten mit entsprechenden Berechtigungen
- Zugriff auf die Region eu-west-1
- Verständnis des Requester-Pays-Abrechnungsmodells
Konfigurationsschritte
1. AWS-Zugangsdaten und Region festlegen
export AWS_ACCESS_KEY_ID="your-access-key"
export AWS_SECRET_ACCESS_KEY="your-secret-key"
export AWS_DEFAULT_REGION="eu-west-1"
2. Bucket-Zugriff überprüfen
aws s3 ls s3://cfg-public-proper-wallaby/ \
--region eu-west-1 \
--request-payer requester
3. Verfügbare Instrumente ermitteln
Listen Sie vor dem Herunterladen alle im Bucket verfügbaren Währungspaare (Instrumente) auf. Verwenden Sie --delimiter / (oder einfaches aws s3 ls), um nur Präfixe der obersten Ebene aufzulisten — vermeiden Sie rekursives Auflisten, da dies jede einzelne Datei durchsuchen würde, nur um die ~20-30 Ordnernamen zu finden.
Schnelle Liste (am einfachsten):
aws s3 ls s3://cfg-public-proper-wallaby/ \
--region eu-west-1 \
--request-payer requester
Saubere Liste (nur Namen, sortiert, in Datei gespeichert):
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 "Gesamtzahl der Instrumente: $(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):
"""Alle Präfixe der obersten Ebene abrufen (Währungspaare)"""
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"Gefundene Instrumente insgesamt: {len(instruments)}\n")
for inst in instruments:
print(inst)
with open('instruments.txt', 'w') as f:
f.write('\n'.join(instruments))
Verwendung:
python3 list_instruments.py
⚠️ Hinweis zu Leistung und Kosten: Die Verwendung von
--delimiter /(bzw.Delimiter='/'in Boto3) ist unerlässlich — sie löst eine schlankeCommonPrefixes-Antwort aus, anstatt jede Datei aufzulisten. Dies kostet nur wenige günstige LIST-Anfragen und dauert Sekunden, im Gegensatz zu einem langsamen, teureren vollständigen Bucket-Scan ohne Trennzeichen.
4. Grundlegender S3-Zugriffsbefehl mit Requester Pays
# Einzelne Datei herunterladen
aws s3 cp s3://cfg-public-proper-wallaby/EURUSD/file.json . \
--region eu-west-1 \
--request-payer requester
5. Batch-Download mit Leistungsoptimierung
# Gesamte Historie für ein bestimmtes Währungspaar herunterladen
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
Parameter zur Leistungsoptimierung
| Parameter | Wert | Vorteil |
|---|---|---|
--max-concurrent-requests |
20-30 | Erhöht parallele Downloads |
--max-bandwidth |
100MB/s | Verhindert Netzwerk-Engpässe |
--no-progress |
- | Reduziert I/O-Overhead |
--region |
eu-west-1 | Reduziert Latenz innerhalb der EU-Region |
6. Erweitertes Download-Skript (Geschwindigkeit & Stabilität)
#!/bin/bash
BUCKET_NAME="cfg-public-proper-wallaby"
REGION="eu-west-1"
CURRENCY_PAIR="${1:-EURUSD}" # Standard: EURUSD, oder als Argument übergeben
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 "Starte S3-Sync für $CURRENCY_PAIR aus $BUCKET_NAME in Region $REGION..." | tee "$LOG_FILE"
for attempt in $(seq 1 $MAX_RETRIES); do
echo "Download-Versuch $attempt von $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 "✓ Download erfolgreich abgeschlossen um $(date)" | tee -a "$LOG_FILE"
echo "Dateien insgesamt: $(find $DESTINATION -type f | wc -l)" | tee -a "$LOG_FILE"
echo "Gesamtgröße: $(du -sh $DESTINATION | cut -f1)" | tee -a "$LOG_FILE"
exit 0
else
echo "✗ Download fehlgeschlagen mit Exit-Code $EXIT_CODE. Erneuter Versuch in 30 Sekunden..." | tee -a "$LOG_FILE"
sleep 30
fi
done
echo "✗ Download nach $MAX_RETRIES Versuchen fehlgeschlagen" | tee -a "$LOG_FILE"
exit 1
Verwendung:
chmod +x download-price-history.sh
./download-price-history.sh EURUSD
./download-price-history.sh GBPUSD
./download-price-history.sh
7. Batch-Download mehrerer Währungspaare
#!/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 "Starte Download für $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 abgeschlossen"
break
else
echo "⚠ $PAIR Versuch $attempt fehlgeschlagen, erneuter Versuch..."
sleep 30
fi
done
done
echo "✓ Alle Downloads abgeschlossen"
8. Verwendung von Python Boto3 (programmatischer Zugriff)
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):
"""Datei mit aktiviertem Requester Pays herunterladen"""
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"✓ Heruntergeladen: {key}")
return True
except Exception as e:
logger.error(f"✗ Fehler beim Herunterladen von {key}: {e}")
return False
def batch_download(bucket, currency_pair, destination, max_workers=10):
"""Alle Dateien für ein Währungspaar parallel herunterladen"""
try:
bucket_obj = s3_resource.Bucket(bucket)
# Ein abschließender Schrägstrich beim Präfix verhindert, dass versehentlich
# ein anderes Instrument getroffen wird, das diesen Text zufällig als Präfix hat
# (z. B. "EURUSD" gegenüber einem hypothetischen "EURUSDT").
prefix = f"{currency_pair}/"
objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
logger.info(f"{len(objects)} Objekte zum Herunterladen für {currency_pair} gefunden")
with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
futures = []
for obj in objects:
if obj.key.endswith('/'): # Verzeichnisse überspringen
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=== Download-Zusammenfassung ===")
logger.info(f"Währungspaar: {currency_pair}")
logger.info(f"Dateien insgesamt: {len(futures)}")
logger.info(f"Erfolgreich: {completed}")
logger.info(f"Fehlgeschlagen: {failed}")
return completed, failed
except Exception as e:
logger.error(f"Fehler beim Batch-Download: {e}")
return 0, len(objects)
# Verwendung - Einzelnes Währungspaar
if __name__ == "__main__":
BUCKET = 'cfg-public-proper-wallaby'
CURRENCY_PAIR = 'EURUSD' # Bei Bedarf ändern
DESTINATION = f'./data/{CURRENCY_PAIR}'
logger.info(f"Starte Download von s3://{BUCKET}/{CURRENCY_PAIR}/")
logger.info(f"Region: eu-west-1")
logger.info(f"Ziel: {DESTINATION}")
completed, failed = batch_download(
bucket=BUCKET,
currency_pair=CURRENCY_PAIR,
destination=DESTINATION,
max_workers=10
)
if failed == 0:
logger.info("\n✓ Alle Dateien erfolgreich heruntergeladen!")
else:
logger.warning(f"\n⚠ {failed} Dateien konnten nicht heruntergeladen werden")
Verwendung:
python3 download_price_history.py
# CURRENCY_PAIR-Variable für andere Paare bearbeiten
Hinweis: Angesichts der bestätigten 26.586 Dateien für EUR/USD listet der obige Aufruf list(bucket_obj.objects.filter(...)) alle über paginierte ListObjectsV2-Aufrufe auf, bevor die Downloads beginnen. Dieser Auflistungsschritt selbst verursacht eigene (geringe) Anfragekosten, die separat von GET-Anfragen abgerechnet werden. Bei sehr großen Präfixen sollten Sie manuelles Paginieren oder das Testen mit Teilabrufen in Betracht ziehen.
9. Alle Währungspaare herunterladen (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):
"""Datei mit aktiviertem Requester Pays herunterladen"""
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"✗ Fehler beim Herunterladen von {key}: {e}")
return False
def download_all_pairs(bucket, destination_base, max_workers=10):
"""Gesamte Historie für alle Währungspaare herunterladen"""
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)} Währungspaare gefunden: {sorted(pairs)}")
total_completed = 0
total_failed = 0
for pair in sorted(pairs):
logger.info(f"\n=== Verarbeite {pair} ===")
# Der abschließende Schrägstrich begrenzt dies exakt auf die Schlüssel dieses Instruments
prefix = f"{pair}/"
pair_objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
logger.info(f"{len(pair_objects)} Objekte für {pair} gefunden")
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)} Dateien heruntergeladen")
logger.info(f"\n=== GESAMTÜBERSICHT ===")
logger.info(f"Insgesamt abgeschlossen: {total_completed}")
logger.info(f"Insgesamt fehlgeschlagen: {total_failed}")
return total_completed, total_failed
except Exception as e:
logger.error(f"Fehler: {e}")
return 0, 0
# Verwendung
if __name__ == "__main__":
BUCKET = 'cfg-public-proper-wallaby'
DESTINATION_BASE = './data'
logger.info(f"Starte vollständigen Archiv-Download von s3://{BUCKET}/")
logger.info(f"Region: eu-west-1")
logger.info(f"Ziel: {DESTINATION_BASE}")
download_all_pairs(
bucket=BUCKET,
destination_base=DESTINATION_BASE,
max_workers=10
)
10. Kurzreferenz-Befehle
# Alle verfügbaren Währungspaare auflisten
aws s3 ls s3://cfg-public-proper-wallaby/ \
--region eu-west-1 \
--request-payer requester
# EURUSD-Historie herunterladen
aws s3 sync s3://cfg-public-proper-wallaby/EURUSD/ ./EURUSD/ \
--region eu-west-1 \
--request-payer requester \
--max-concurrent-requests 25
# GBPUSD-Historie herunterladen
aws s3 sync s3://cfg-public-proper-wallaby/GBPUSD/ ./GBPUSD/ \
--region eu-west-1 \
--request-payer requester \
--max-concurrent-requests 25
# Dateien in EURUSD zählen
aws s3 ls s3://cfg-public-proper-wallaby/EURUSD/ \
--region eu-west-1 \
--request-payer requester \
--recursive | wc -l
# Gesamtgröße von EURUSD ermitteln
aws s3 ls s3://cfg-public-proper-wallaby/EURUSD/ \
--region eu-west-1 \
--request-payer requester \
--recursive --summarize | grep "Total Size"
Wichtige Beweggründe
Geschwindigkeitssteigerung ⚡
- Parallelverarbeitung: 20-25 gleichzeitige Verbindungen maximieren den Durchsatz
- EU-West-1-Optimierung: keine regionsübergreifende Latenz
- Adaptive Wiederholungsstrategie: intelligente Fehlerbehebung
- Connection Pooling: minimiert Overhead zwischen Anfragen
- Bandbreitenkontrolle: verhindert Netzwerküberlastung und Stabilitätsprobleme
Stabilitätsverbesserung
- Wiederholungsmechanismus: 3 Versuche mit 30-Sekunden-Intervallen
- Fehlerprotokollierung: detaillierte Logs zur Fehlerbehebung
- Verbindungs-Timeout: 5 Sekunden für Verbindungsaufbau, 60 Sekunden für Lesevorgänge
- Adaptiver Modus: passt die Wiederholungsstrategie je nach Fehlertyp an
- Integritätsprüfungen: validieren heruntergeladene Dateien
- Graceful Degradation: läuft bei Teilfehlern weiter
Kostenschätzung
Preisstruktur
AWS-S3-GET-Anfragepreise (Requester Pays, eu-west-1):
- $0,0004 pro 1.000 Anfragen
- $0,02 pro GB Datenübertragung (eu-west-1-Ausgangstarif)
Die Dekodierungs- und Pipeline-Schritte in den Abschnitten „Dekodierung von
.bi5-Dateien" und „Vollständige Pipeline" laufen vollständig auf Ihrem lokalen Rechner — sie verursachen keine zusätzlichen AWS-Kosten über die hier genannten Download-Kosten hinaus.
Schätzung des Gesamtarchivs
⚠️ Gesamtarchiv: ~400 GB, ~20.000.000 Dateien
Anfragekosten = (20.000.000 / 1.000) × $0,0004 = $8,00
Übertragungskosten = 400 × $0,02 = $8,00
─────────────────────────────────────────────────────────
GESAMT = $16,00
Durchschnittliche Dateigröße im gesamten Archiv: 400 GB / 20.000.000 Dateien ≈ 20,97 KB/Datei
Praxisbeispiel: EUR/USD (verifizierte tatsächliche Daten)
Bestätigt über aws s3 ls --summarize:
Objekte insgesamt: 26.586
Gesamtgröße: 2,6 GB
Durchschnittliche Dateigröße: 2,6 GB / 26.586 ≈ 100,1 KB/Datei
Anfragekosten:
Anzahl der Anfragen: 26.586
Kosten pro 1.000 Anfragen: $0,0004
Gesamte Anfragekosten = (26.586 / 1.000) × $0,0004
= 26,586 × $0,0004
≈ $0,0106
Datenübertragungskosten:
Datengröße: 2,6 GB
Kosten pro GB: $0,02
Gesamte Übertragungskosten = 2,6 × $0,02 = $0,052
Gesamtkosten für EUR/USD-Download (verifiziert):
Anfragekosten: $0,0106
Datenübertragungskosten: $0,052
────────────────────────────
GESAMT: ~$0,06
Kostenvergleichstabelle
| Szenario | Dateien | Größe | Anfragekosten | Übertragungskosten | Gesamtkosten |
|---|---|---|---|---|---|
| Gesamtarchiv | 20M | 400 GB | $8,00 | $8,00 | $16,00 |
| EUR/USD (verifiziert, tatsächlich) | 26.586 | 2,6 GB | $0,0106 | $0,052 | ~$0,06 |
| 5 Paare (falls ähnlich wie EURUSD) | ~133.000 | ~13 GB | $0,053 | $0,26 | ~$0,31 |
| 10 Paare (falls ähnlich wie EURUSD) | ~266.000 | ~26 GB | $0,106 | $0,52 | ~$0,63 |
Die Schätzungen für „5 Paare" und „10 Paare" gehen von einer Datei-/Größenanzahl ähnlich wie bei EUR/USD aus. Tatsächliche Kosten variieren je nach Paar — Hauptwährungspaare könnten mehr Historie/Dateien als EUR/USD aufweisen, exotische Paare vermutlich weniger.
Tipps zur Kostenoptimierung
- ✓ Laden Sie nur benötigte Währungspaare herunter (nicht das gesamte Archiv)
- ✓ Überprüfen Sie immer die tatsächliche Größe/Dateianzahl pro Paar mit
--summarize— die tatsächliche durchschnittliche Dateigröße von EUR/USD (~100 KB) ist etwa 5-mal so groß wie der archivweite Durchschnitt (~21 KB), was zeigt, wie irreführend Durchschnittswerte sein können - ✓ Bündeln Sie Downloads, um den Anfrage-Overhead zu minimieren
- ✓ Verwenden Sie das Flag
--delete, um erneutes Herunterladen bereits vorhandener Dateien zu vermeiden - ✓ Bei ~$0,06 pro Paar (verifiziert an EUR/USD) bleibt das einzelne Herunterladen von Dutzenden Paaren weiterhin nur ein kleiner Bruchteil der $16-Gesamtarchivkosten
- ✓ Verwenden Sie Abschnitt 3 (Verfügbare Instrumente ermitteln), um genaue Paarnamen zu bestätigen, bevor Sie Batch-Skripte ausführen
Empfohlene Maßnahme vor Massendownloads
aws s3 ls s3://cfg-public-proper-wallaby/<PAIR>/ \
--region eu-west-1 \
--request-payer requester \
--recursive --summarize | tail -3
Dies liefert die genauen Werte für Total Objects und Total Size und ermöglicht die Berechnung genauer Kosten mittels:
Anfragekosten = (Dateien / 1.000) × $0,0004
Übertragungskosten = Größe_in_GB × $0,02
Dekodierung von .bi5-Dateien in lesbare Tick-Daten
Die .bi5-Dateien von Dukascopy sind LZMA-komprimierte binäre Tick-Dateien. Bei der aktuellen Bucket-Struktur stellt jede Datei einen vollen Tag an Tick-Daten pro Instrument dar.
Dateistruktur
Komprimierung: Roher LZMA-Stream (kein .xz-Containerformat)
Pfadkonvention (täglich):
SYMBOL/YEAR/MONTH/DAY_ticks.bi5
Beispiel: EURUSD/2024/00/15_ticks.bi5 → EUR/USD-Ticks für den 15. Januar 2024
⚠️ Der Monat ist nullbasiert (Januar =
00, Dezember =11).✅ Keine leeren Dateien: Eine fehlende Datei für einen bestimmten Tag bedeutet, dass an diesem Tag keine Ticks erfasst wurden (Wochenenden, Feiertage). Behandeln Sie fehlende S3-Schlüssel/
FileNotFoundErrorals „keine Daten" — nicht als Fehler.
Format des dekomprimierten Datensatzes
Jeder Tick ist ein fester binärer Datensatz von 20 Byte, Big-Endian:
| Byte | Feld | Typ | Hinweise |
|---|---|---|---|
| 0–3 | Zeitstempel | uint32 | Millisekunden seit Tagesbeginn (UTC) |
| 4–7 | Ask-Preis | uint32 | Erfordert Skalierung nach Punktwert |
| 8–11 | Bid-Preis | uint32 | Erfordert Skalierung nach Punktwert |
| 12–15 | Ask-Volumen | float32 | In Millionen Einheiten der Basiswährung |
| 16–19 | Bid-Volumen | float32 | In Millionen Einheiten der Basiswährung |
Preisskalierung („Punktwert")
| Instrumententyp | Punktwert | Beispiel |
|---|---|---|
| Die meisten FX-Paare | 100.000 | EUR/USD, GBP/USD |
| JPY-Paare | 1.000 | USD/JPY, EUR/JPY |
| Indizes/Rohstoffe | Variiert | Pro Instrument überprüfen |
⚠️ Die untenstehenden Hilfsfunktionen unterscheiden nur zwischen „JPY-Paar" und „allem anderen (100.000)". Bestätigen Sie bei Indizes, Rohstoffen oder jedem anderen Nicht-FX-Instrument den korrekten Punktwert beim Datenanbieter, bevor Sie dekodieren — das stillschweigende Anwenden von 100.000 auf ein Nicht-FX-Instrument führt zu falschen Preisen, ohne dass ein Fehler ausgelöst wird. Die aktualisierte
get_point_value()-Funktion in dieser Anleitung gibt jetzt eine Warnung aus, wenn für ein nicht erkanntes, nicht-JPY-Instrument auf den Standardwert zurückgegriffen wird, damit dies nicht unbemerkt bleibt.
Python-Decoder (Einzelner Tag)
import lzma
import struct
from datetime import datetime, timedelta
def decode_bi5_daily(filepath, day_start, point_value=100000):
"""
Dekodiert eine tägliche .bi5-Tick-Datei.
:param filepath: Pfad zur .bi5-Datei
:param day_start: datetime, das 00:00 UTC dieses Tages darstellt
:param point_value: 100000 für die meisten Paare, 1000 für JPY-Paare
:return: Liste von Tick-Dictionaries
"""
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
# Beispielverwendung
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"Ticks insgesamt für den Tag: {len(ticks)}")
for tick in ticks[:5]:
print(tick)
Batch-Dekodierungsskript (kompletter Instrumentordner → 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'}
# Bekannte Nicht-FX-Instrumente (Indizes/Rohstoffe) mit bestätigtem Punktwert.
# Erweitern Sie diese Liste, sobald Sie weitere Instrumente beim Datenanbieter bestätigt haben.
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}' sieht nicht wie ein standardmäßiges 6-Buchstaben-FX-Paar aus und hat "
f"keinen bestätigten Punktwert — es wird der Standardwert 100000 verwendet. Überprüfen Sie dies, bevor Sie dem Ergebnis vertrauen."
)
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):
"""
Dekodiert alle täglichen .bi5-Dateien für ein Instrument in eine einzelne CSV-Datei.
Erwartet die Ordnerstruktur: 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"✗ Fehler beim Dekodieren von {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=== Dekodierungs-Zusammenfassung ===")
print(f"Instrument: {instrument_name}")
print(f"Verarbeitete Dateien: {file_count}")
print(f"Fehler: {error_count}")
print(f"Ticks insgesamt: {len(all_ticks)}")
print(f"Ausgabe: {output_csv}")
if __name__ == "__main__":
batch_decode_instrument('./data/EURUSD', './EURUSD_ticks.csv')
Verwendung:
python3 decode_bi5_batch.py
Häufige Fallstricke
- ⚠️ Der Monat ist nullbasiert im S3-Pfad, aber nicht in der Datumsberechnung — explizit umrechnen.
- ⚠️ Falscher Punktwert: Eine Diskrepanz zwischen 100.000 und 1.000 korrumpiert Preise stillschweigend. Nicht-FX-Instrumente (Indizes/Rohstoffe) werden von dieser Regel überhaupt nicht abgedeckt — bestätigen Sie deren Punktwert explizit (siehe
OTHER_POINT_VALUESoben). - ⚠️ Fehlende Dateien ≠ Fehler: Als „an diesem Tag keine Ticks" behandeln, nicht als Dekodierungsfehler.
- ⚠️ Roher LZMA-Stream: Manche Bibliotheken erfordern einen expliziten „Raw"-Modus — Standard-
.xz-Dekomprimierer schlagen fehl. - ⚠️ Ältere stündliche Dateien: Ältere Dateien könnten Millisekunden seit Stundenbeginn statt Millisekunden seit Tagesbeginn verwenden. Keiner der Decoder in dieser Anleitung erkennt dies automatisch — beide gehen durchgehend vom aktuellen täglichen Format (
SYMBOL/YEAR/MONTH/DAY_ticks.bi5) aus. Wenn Sie mit einem Instrument oder Datumsbereich arbeiten, der alt genug ist, um dem täglichen Format vorauszugehen, überprüfen Sie die tatsächliche Zeitbasis der Datei vor dem Dekodieren, anstatt anzunehmen, dass sie dem hier beschriebenen täglichen Format entspricht.
End-to-End-Pipeline (Download → Dekodierung → Export)
Kombiniert Download, Dekodierung und Export in einem automatisierten Skript — mit optionaler Datumsfilterung und Bereinigung.
Abhängigkeiten
pip install boto3 pandas pyarrow
pyarrowwird nur für den Parquet-Export benötigt (empfohlen für große Datensätze).
Vollständiges Pipeline-Skript
#!/usr/bin/env python3
"""
End-to-End-Pipeline: Download -> Dekodierung -> Export
Tägliche .bi5-Tick-Dateien im Dukascopy-Stil aus 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}' sieht nicht wie ein standardmäßiges 6-Buchstaben-FX-Paar aus — "
f"es wird der Standardwert 100000 verwendet. Überprüfen Sie dies, bevor Sie dem Ergebnis vertrauen."
)
return 100000
# Schritt 1: Download
def download_instrument(instrument, destination, start_date=None, end_date=None):
"""Lädt alle .bi5-Dateien für ein Instrument herunter, optional gefiltert nach Datumsbereich."""
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"Objekte für {instrument} werden aufgelistet...")
objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
logger.info(f"{len(objects)} Objekte gefunden")
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"Datum konnte nicht aus Schlüssel geparst werden: {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"✗ Herunterladen von {obj.key} fehlgeschlagen: {e}")
failed += 1
logger.info(f"Download abgeschlossen: {downloaded} heruntergeladen, {skipped} übersprungen (Datumsfilter), {failed} fehlgeschlagen")
return downloaded, failed
# Schritt 2: Dekodierung
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):
"""Dekodiert alle täglichen .bi5-Dateien in einem Ordner in eine Liste von Tick-Tupeln."""
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"✗ Fehler beim Dekodieren von {bi5_file}: {e}")
error_count += 1
all_ticks.sort(key=lambda t: t[0])
logger.info(f"{file_count} Dateien dekodiert ({error_count} Fehler), {len(all_ticks)} Ticks insgesamt")
return all_ticks
# Schritt 3: Export
def export_ticks(ticks, output_path, output_format='csv'):
"""Exportiert dekodierte Ticks mit pandas in CSV oder Parquet."""
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 nach {output_path} ({output_format}) exportiert")
return df
# Pipeline-Orchestrierung
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 gestartet für {instrument} ===")
if start_date or end_date:
logger.info(f"Datumsbereich: {start_date or 'frühester'} bis {end_date or 'spätester'}")
downloaded, failed = download_instrument(instrument, raw_dir, start_date, end_date)
if failed:
logger.warning(f"{failed} Datei(en) konnten für {instrument} nicht heruntergeladen werden — es wird mit den {downloaded} erfolgreichen fortgefahren")
if downloaded == 0:
logger.warning(f"Keine Dateien für {instrument} heruntergeladen — Dekodierung/Export wird übersprungen")
return
ticks = decode_instrument_folder(raw_dir, instrument)
if not ticks:
logger.warning(f"Keine Ticks für {instrument} dekodiert")
return
export_ticks(ticks, output_path, output_format)
if cleanup:
logger.info(f"Bereinige rohe .bi5-Dateien unter {raw_dir}...")
shutil.rmtree(raw_dir, ignore_errors=True)
logger.info("✓ Bereinigung abgeschlossen")
logger.info(f"=== Pipeline abgeschlossen für {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
)
Verwendungsbeispiele
# EURUSD herunterladen, dekodieren und nach Parquet exportieren (Standard)
python3 pipeline.py EURUSD
# Stattdessen nach CSV exportieren
python3 pipeline.py EURUSD --format csv
# Auf einen bestimmten Datumsbereich filtern
python3 pipeline.py EURUSD --start-date 2024-01-01 --end-date 2024-03-31
# Rohe .bi5-Dateien nach dem Export bereinigen
python3 pipeline.py EURUSD --cleanup
# Benutzerdefiniertes Ausgabeverzeichnis
python3 pipeline.py GBPUSD --output-dir ./exports --cleanup
Batch-Pipeline für mehrere Instrumente
#!/bin/bash
# run_all_pipelines.sh
INSTRUMENTS=("EURUSD" "GBPUSD" "USDJPY" "AUDUSD")
for INSTRUMENT in "${INSTRUMENTS[@]}"; do
echo "=== Verarbeite $INSTRUMENT ==="
python3 pipeline.py "$INSTRUMENT" --format parquet --cleanup
done
echo "✓ Alle Pipelines abgeschlossen"
Best Practices ✅
Download & Kostenkontrolle
- ✓ Immer die Flags
--region eu-west-1und--request-payer requestereinbeziehen - ✓ Vor Massendownloads die tatsächliche Größe/Dateianzahl pro Paar mit
--summarizeüberprüfen — Durchschnittswerte können irreführend sein (z. B. ist die reale durchschnittliche Dateigröße von EUR/USD ~5-mal so groß wie der archivweite Durchschnitt) - ✓
aws s3 sync(nichtcp) für mehrere Dateien verwenden — es werden nur geänderte Daten übertragen - ✓ Das Flag
--deleteverwenden, um erneutes Herunterladen bereits vorhandener Dateien zu vermeiden - ✓ Zuerst die Instrumentenermittlung ausführen, um exakte, gültige Paarnamen zu bestätigen — vermeidet fehlgeschlagene Synchronisierungen durch Tippfehler
- ✓ Log-Dateien für Audit-Trails aufbewahren
- ✓ CloudWatch-Alarme für fehlgeschlagene Downloads einrichten
- ✓ Kosten neu berechnen, wenn sich Bucket-Größe/Dateianzahl wesentlich ändert
- ✓ Nach Möglichkeit das Herunterladen bestimmter Paare gegenüber dem gesamten Archiv bevorzugen
Dekodierung & Datenintegrität
- ✓ Vor der Dekodierung immer den korrekten Punktwert (100.000 vs. 1.000) pro Instrument bestätigen — falsche Skalierung korrumpiert Preise stillschweigend
- ✓ Fehlende tägliche
.bi5-Dateien als „an diesem Tag keine Ticks" behandeln (Wochenenden/Feiertage), nicht als Fehler - ✓ Niemals ältere stündliche
.bi5-Dateien mit neuen täglichen Dateien in derselben Pipeline mischen, ohne zuerst das Format zu erkennen - ✓ Den LZMA-Dekomprimierungsmodus überprüfen — manche Bibliotheken erfordern einen expliziten „Raw"-Modus (ohne Container-Header)
- ✓ Dekodierte Ticks vor dem Export chronologisch sortieren, falls es Inkonsistenzen in der Dateisystem-Reihenfolge gibt
Pipeline & Automatisierung
- ✓
--cleanupverwenden, um rohe.bi5-Dateien nach der Dekodierung zu löschen — sie können bei Bedarf jederzeit erneut von S3 heruntergeladen werden - ✓ Parquet gegenüber CSV bevorzugen für Tick-Level-Daten — typischerweise 5-10-mal kleiner und schneller abfragbar
- ✓
--start-date/--end-date-Filter beim Testen verwenden, um unnötige S3-Kosten zu vermeiden - ✓ Mit der Instrumentenermittlung kombinieren, um automatisch alle gültigen Paare zu durchlaufen, statt Listen fest zu codieren
- ✓ Jeden Pipeline-Lauf protokollieren (Anzahl heruntergeladen/dekodiert/exportiert) zur Nachvollziehbarkeit
Allgemein
- ✓ Downloads außerhalb der Stoßzeiten durchführen, um Konkurrenz zu reduzieren (nicht Kosten)
- ✓ Prüfsummen zur Überprüfung der Datenintegrität implementieren, wo dies kritisch ist
- ✓ Die insgesamt übertragene Datenmenge im Zeitverlauf überwachen, um kumulative Kosten zu kontrollieren