Dati storici sui prezzi
Panoramica
Questa guida consente di accedere ai file di dati storici sui prezzi dal bucket S3 cfg-public-proper-wallaby (eu-west-1), configurato con il criterio Requester Pays, garantendo prestazioni di download ottimali e stabilità del sistema.
Prerequisiti
- AWS CLI installato e configurato
- Credenziali AWS valide con i permessi appropriati
- Accesso alla regione eu-west-1
- Comprensione del modello di fatturazione requester pays
Passaggi di configurazione
1. Impostare le credenziali AWS e la regione
export AWS_ACCESS_KEY_ID="your-access-key"
export AWS_SECRET_ACCESS_KEY="your-secret-key"
export AWS_DEFAULT_REGION="eu-west-1"
2. Verificare l'accesso al bucket
aws s3 ls s3://cfg-public-proper-wallaby/ \
--region eu-west-1 \
--request-payer requester
3. Individuare gli strumenti disponibili
Prima di scaricare, elenca tutte le coppie di valute (strumenti) disponibili nel bucket. Usa --delimiter / (oppure il semplice aws s3 ls) per elencare solo i prefissi di primo livello — evita l'elenco ricorsivo, che scansionerebbe ogni singolo file solo per trovare i ~20-30 nomi di cartelle.
Elenco rapido (il più semplice):
aws s3 ls s3://cfg-public-proper-wallaby/ \
--region eu-west-1 \
--request-payer requester
Elenco pulito (solo nomi, ordinati, salvati su file):
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 "Totale strumenti: $(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):
"""Ottiene tutti i prefissi di primo livello (coppie di valute)"""
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"Totale strumenti trovati: {len(instruments)}\n")
for inst in instruments:
print(inst)
with open('instruments.txt', 'w') as f:
f.write('\n'.join(instruments))
Utilizzo:
python3 list_instruments.py
⚠️ Nota su prestazioni e costi: L'uso di
--delimiter /(oDelimiter='/'in Boto3) è essenziale — genera una risposta leggeraCommonPrefixesinvece di elencare ogni singolo file. Costa solo poche richieste LIST economiche e si completa in pochi secondi, contro una scansione completa del bucket lenta e più costosa senza il delimitatore.
4. Comando di accesso S3 di base con Requester Pays
# Scaricare un singolo file
aws s3 cp s3://cfg-public-proper-wallaby/EURUSD/file.json . \
--region eu-west-1 \
--request-payer requester
5. Download in blocco con ottimizzazione delle prestazioni
# Scaricare l'intera storia per una specifica coppia di valute
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
Parametri di ottimizzazione delle prestazioni
| Parametro | Valore | Vantaggio |
|---|---|---|
--max-concurrent-requests |
20-30 | Aumenta i download paralleli |
--max-bandwidth |
100MB/s | Previene i colli di bottiglia di rete |
--no-progress |
- | Riduce il sovraccarico di I/O |
--region |
eu-west-1 | Riduce la latenza all'interno della regione UE |
6. Script di download avanzato (velocità e stabilità)
#!/bin/bash
BUCKET_NAME="cfg-public-proper-wallaby"
REGION="eu-west-1"
CURRENCY_PAIR="${1:-EURUSD}" # Predefinito EURUSD, oppure passato come argomento
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 "Avvio sincronizzazione S3 per $CURRENCY_PAIR da $BUCKET_NAME nella regione $REGION..." | tee "$LOG_FILE"
for attempt in $(seq 1 $MAX_RETRIES); do
echo "Tentativo di download $attempt di $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 completato con successo il $(date)" | tee -a "$LOG_FILE"
echo "File totali: $(find $DESTINATION -type f | wc -l)" | tee -a "$LOG_FILE"
echo "Dimensione totale: $(du -sh $DESTINATION | cut -f1)" | tee -a "$LOG_FILE"
exit 0
else
echo "✗ Download fallito con codice di uscita $EXIT_CODE. Nuovo tentativo tra 30 secondi..." | tee -a "$LOG_FILE"
sleep 30
fi
done
echo "✗ Download fallito dopo $MAX_RETRIES tentativi" | tee -a "$LOG_FILE"
exit 1
Utilizzo:
chmod +x download-price-history.sh
./download-price-history.sh EURUSD
./download-price-history.sh GBPUSD
./download-price-history.sh
7. Download in blocco di più coppie di valute
#!/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 "Avvio download per $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 completato"
break
else
echo "⚠ $PAIR tentativo $attempt fallito, nuovo tentativo..."
sleep 30
fi
done
done
echo "✓ Tutti i download completati"
8. Utilizzo di Python Boto3 (accesso programmatico)
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):
"""Scarica il file con Requester Pays abilitato"""
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"✓ Scaricato: {key}")
return True
except Exception as e:
logger.error(f"✗ Errore durante il download di {key}: {e}")
return False
def batch_download(bucket, currency_pair, destination, max_workers=10):
"""Scarica tutti i file per una coppia di valute in parallelo"""
try:
bucket_obj = s3_resource.Bucket(bucket)
# Una barra finale nel prefisso evita di corrispondere accidentalmente
# a un altro strumento che ha questo testo come prefisso
# (ad es. "EURUSD" rispetto a un ipotetico "EURUSDT").
prefix = f"{currency_pair}/"
objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
logger.info(f"Trovati {len(objects)} oggetti da scaricare per {currency_pair}")
with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
futures = []
for obj in objects:
if obj.key.endswith('/'): # Saltare le directory
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=== Riepilogo download ===")
logger.info(f"Coppia di valute: {currency_pair}")
logger.info(f"File totali: {len(futures)}")
logger.info(f"Riusciti: {completed}")
logger.info(f"Falliti: {failed}")
return completed, failed
except Exception as e:
logger.error(f"Errore nel download in blocco: {e}")
return 0, len(objects)
# Utilizzo - Singola coppia di valute
if __name__ == "__main__":
BUCKET = 'cfg-public-proper-wallaby'
CURRENCY_PAIR = 'EURUSD' # Modificare secondo necessità
DESTINATION = f'./data/{CURRENCY_PAIR}'
logger.info(f"Avvio download da s3://{BUCKET}/{CURRENCY_PAIR}/")
logger.info(f"Regione: eu-west-1")
logger.info(f"Destinazione: {DESTINATION}")
completed, failed = batch_download(
bucket=BUCKET,
currency_pair=CURRENCY_PAIR,
destination=DESTINATION,
max_workers=10
)
if failed == 0:
logger.info("\n✓ Tutti i file scaricati con successo!")
else:
logger.warning(f"\n⚠ {failed} file non sono stati scaricati")
Utilizzo:
python3 download_price_history.py
# Modificare la variabile CURRENCY_PAIR per coppie diverse
Nota: Dato che per EUR/USD sono stati confermati 26.586 file, la chiamata list(bucket_obj.objects.filter(...)) sopra elencherà tutti questi file tramite chiamate paginate ListObjectsV2 prima che inizino i download. Questo passaggio di elenco comporta di per sé un proprio (piccolo) costo di richiesta, fatturato separatamente dalle richieste GET. Per prefissi molto grandi, valuta la paginazione manuale o il test con recuperi parziali.
9. Scaricare tutte le coppie di valute (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):
"""Scarica il file con Requester Pays abilitato"""
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"✗ Errore durante il download di {key}: {e}")
return False
def download_all_pairs(bucket, destination_base, max_workers=10):
"""Scarica l'intera storia per tutte le coppie di valute"""
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"Trovate {len(pairs)} coppie di valute: {sorted(pairs)}")
total_completed = 0
total_failed = 0
for pair in sorted(pairs):
logger.info(f"\n=== Elaborazione {pair} ===")
# La barra finale limita la ricerca esattamente alle chiavi di questo strumento
prefix = f"{pair}/"
pair_objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
logger.info(f"Trovati {len(pair_objects)} oggetti per {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)} file scaricati")
logger.info(f"\n=== RIEPILOGO FINALE ===")
logger.info(f"Totale completati: {total_completed}")
logger.info(f"Totale falliti: {total_failed}")
return total_completed, total_failed
except Exception as e:
logger.error(f"Errore: {e}")
return 0, 0
# Utilizzo
if __name__ == "__main__":
BUCKET = 'cfg-public-proper-wallaby'
DESTINATION_BASE = './data'
logger.info(f"Avvio download dell'archivio completo da s3://{BUCKET}/")
logger.info(f"Regione: eu-west-1")
logger.info(f"Destinazione: {DESTINATION_BASE}")
download_all_pairs(
bucket=BUCKET,
destination_base=DESTINATION_BASE,
max_workers=10
)
10. Comandi di riferimento rapido
# Elencare tutte le coppie di valute disponibili
aws s3 ls s3://cfg-public-proper-wallaby/ \
--region eu-west-1 \
--request-payer requester
# Scaricare la storia di EURUSD
aws s3 sync s3://cfg-public-proper-wallaby/EURUSD/ ./EURUSD/ \
--region eu-west-1 \
--request-payer requester \
--max-concurrent-requests 25
# Scaricare la storia di GBPUSD
aws s3 sync s3://cfg-public-proper-wallaby/GBPUSD/ ./GBPUSD/ \
--region eu-west-1 \
--request-payer requester \
--max-concurrent-requests 25
# Contare i file in EURUSD
aws s3 ls s3://cfg-public-proper-wallaby/EURUSD/ \
--region eu-west-1 \
--request-payer requester \
--recursive | wc -l
# Ottenere la dimensione totale di EURUSD
aws s3 ls s3://cfg-public-proper-wallaby/EURUSD/ \
--region eu-west-1 \
--request-payer requester \
--recursive --summarize | grep "Total Size"
Punti chiave di motivazione
Miglioramento della velocità ⚡
- Elaborazione parallela: 20-25 connessioni simultanee massimizzano il throughput
- Ottimizzazione EU-West-1: nessuna latenza tra regioni diverse
- Strategia di ripetizione adattiva: recupero intelligente dai fallimenti
- Connection pooling: minimizza il sovraccarico tra le richieste
- Controllo della banda: previene la saturazione della rete e problemi di stabilità
Miglioramento della stabilità
- Meccanismo di ripetizione: 3 tentativi con intervalli di 30 secondi
- Registrazione degli errori: log dettagliati per la risoluzione dei problemi
- Timeout di connessione: 5 secondi per la connessione, 60 secondi per la lettura
- Modalità adattiva: adatta la strategia di ripetizione in base ai tipi di errore
- Controlli di integrità: convalidano i download dei file
- Degradazione controllata: continua in caso di fallimenti parziali
Stima dei costi
Struttura dei prezzi
Prezzi delle richieste GET di AWS S3 (Requester Pays, eu-west-1):
- $0,0004 ogni 1.000 richieste
- $0,02 per GB di trasferimento dati (tariffa di uscita eu-west-1)
I passaggi di decodifica e pipeline delle sezioni "Decodifica dei file
.bi5" e "Pipeline end-to-end" vengono eseguiti interamente sul tuo computer locale — non comportano costi AWS aggiuntivi oltre ai costi di download indicati qui.
Stima dell'archivio completo
⚠️ Archivio completo: ~400 GB, ~20.000.000 file
Costo delle richieste = (20.000.000 / 1.000) × $0,0004 = $8,00
Costo del trasferimento = 400 × $0,02 = $8,00
─────────────────────────────────────────────────────────
TOTALE = $16,00
Dimensione media del file per l'intero archivio: 400 GB / 20.000.000 file ≈ 20,97 KB/file
Esempio reale: EUR/USD (dati effettivi verificati)
Confermato tramite aws s3 ls --summarize:
Oggetti totali: 26.586
Dimensione totale: 2,6 GB
Dimensione media del file: 2,6 GB / 26.586 ≈ 100,1 KB/file
Costi delle richieste:
Numero di richieste: 26.586
Costo ogni 1.000 richieste: $0,0004
Costo totale delle richieste = (26.586 / 1.000) × $0,0004
= 26,586 × $0,0004
≈ $0,0106
Costi di trasferimento dati:
Dimensione dei dati: 2,6 GB
Costo per GB: $0,02
Costo totale del trasferimento = 2,6 × $0,02 = $0,052
Costo totale del download EUR/USD (verificato):
Costi delle richieste: $0,0106
Costi di trasferimento dati: $0,052
────────────────────────────
TOTALE: ~$0,06
Tabella comparativa dei costi
| Scenario | File | Dimensione | Costo richieste | Costo trasferimento | Costo totale |
|---|---|---|---|---|---|
| Archivio completo | 20M | 400 GB | $8,00 | $8,00 | $16,00 |
| EUR/USD (verificato, reale) | 26.586 | 2,6 GB | $0,0106 | $0,052 | ~$0,06 |
| 5 coppie (se simili a EURUSD) | ~133.000 | ~13 GB | $0,053 | $0,26 | ~$0,31 |
| 10 coppie (se simili a EURUSD) | ~266.000 | ~26 GB | $0,106 | $0,52 | ~$0,63 |
Le stime per "5 coppie" e "10 coppie" presuppongono un numero di file/dimensioni simili a EUR/USD. I costi effettivi varieranno in base alla coppia — le coppie principali potrebbero avere più storia/file rispetto a EUR/USD, quelle esotiche probabilmente meno.
Suggerimenti per l'ottimizzazione dei costi
- ✓ Scarica solo le coppie di valute necessarie (non l'intero archivio)
- ✓ Verifica sempre la dimensione/numero di file effettivi per coppia con
--summarize— la dimensione media effettiva del file di EUR/USD (~100 KB) è circa 5 volte quella media dell'intero archivio (~21 KB), a dimostrazione di quanto le medie possano essere fuorvianti - ✓ Raggruppa i download per minimizzare il sovraccarico delle richieste
- ✓ Usa il flag
--deleteper evitare di riscaricare file già esistenti - ✓ A ~$0,06 per coppia (verificato su EUR/USD), scaricare individualmente anche decine di coppie resta comunque solo una piccola frazione del costo dell'archivio completo, pari a $16
- ✓ Usa la Sezione 3 (Individuare gli strumenti disponibili) per confermare i nomi esatti delle coppie prima di eseguire gli script in blocco
Azione consigliata prima dei download in blocco
aws s3 ls s3://cfg-public-proper-wallaby/<PAIR>/ \
--region eu-west-1 \
--request-payer requester \
--recursive --summarize | tail -3
Questo restituisce i valori esatti di Total Objects e Total Size, permettendo di calcolare un costo accurato tramite:
Costo delle richieste = (file / 1.000) × $0,0004
Costo del trasferimento = dimensione_in_GB × $0,02
Decodifica dei file .bi5 in dati tick leggibili
I file .bi5 di Dukascopy sono file binari di tick compressi con LZMA. Con l'attuale struttura del bucket, ogni file rappresenta un giorno intero di dati tick per strumento.
Struttura dei file
Compressione: flusso LZMA grezzo (non formato contenitore .xz)
Convenzione dei percorsi (giornaliera):
SYMBOL/YEAR/MONTH/DAY_ticks.bi5
Esempio: EURUSD/2024/00/15_ticks.bi5 → tick EUR/USD del 15 gennaio 2024
⚠️ Il mese è indicizzato da zero (gennaio =
00, dicembre =11).✅ Nessun file vuoto: un file mancante per un determinato giorno significa che quel giorno non sono stati registrati tick (fine settimana, festività). Tratta le chiavi S3 mancanti /
FileNotFoundErrorcome "nessun dato", non come un errore.
Formato del record decompresso
Ogni tick è un record binario fisso di 20 byte, big-endian:
| Byte | Campo | Tipo | Note |
|---|---|---|---|
| 0–3 | Timestamp | uint32 | Millisecondi dall'inizio del giorno (UTC) |
| 4–7 | Prezzo Ask | uint32 | Richiede la scalatura tramite il valore del punto |
| 8–11 | Prezzo Bid | uint32 | Richiede la scalatura tramite il valore del punto |
| 12–15 | Volume Ask | float32 | In milioni di unità della valuta base |
| 16–19 | Volume Bid | float32 | In milioni di unità della valuta base |
Scalatura del prezzo ("valore del punto")
| Tipo di strumento | Valore del punto | Esempio |
|---|---|---|
| La maggior parte delle coppie FX | 100.000 | EUR/USD, GBP/USD |
| Coppie JPY | 1.000 | USD/JPY, EUR/JPY |
| Indici/materie prime | Variabile | Verificare per ogni strumento |
⚠️ Le funzioni di supporto sottostanti distinguono solo tra "coppia JPY" e "tutto il resto (100.000)". Per indici, materie prime o qualsiasi strumento non-FX, conferma il valore del punto corretto con il fornitore dei dati prima della decodifica — applicare silenziosamente 100.000 a uno strumento non-FX produrrà prezzi errati senza generare alcun errore. La funzione
get_point_value()aggiornata in questa guida ora genera un avviso quando ricorre al valore predefinito per uno strumento non riconosciuto e non-JPY, affinché ciò non passi inosservato.
Decodificatore Python (giorno singolo)
import lzma
import struct
from datetime import datetime, timedelta
def decode_bi5_daily(filepath, day_start, point_value=100000):
"""
Decodifica un file tick .bi5 giornaliero.
:param filepath: Percorso del file .bi5
:param day_start: datetime che rappresenta le 00:00 UTC di quel giorno
:param point_value: 100000 per la maggior parte delle coppie, 1000 per le coppie JPY
:return: Elenco di dizionari di tick
"""
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
# Esempio di utilizzo
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"Totale tick per il giorno: {len(ticks)}")
for tick in ticks[:5]:
print(tick)
Script di decodifica in blocco (cartella completa dello strumento → 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'}
# Strumenti non-FX noti (indici/materie prime) con valore del punto confermato.
# Espandi questo elenco man mano che confermi altri strumenti con il fornitore dei dati.
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}' non sembra una coppia FX standard a 6 lettere e non ha "
f"un valore del punto confermato — verrà usato il valore predefinito 100000. Verifica questo prima di fidarti del risultato."
)
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):
"""
Decodifica tutti i file .bi5 giornalieri di uno strumento in un unico CSV.
Prevede la struttura di cartelle: 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"✗ Errore durante la decodifica di {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=== Riepilogo decodifica ===")
print(f"Strumento: {instrument_name}")
print(f"File elaborati: {file_count}")
print(f"Errori: {error_count}")
print(f"Totale tick: {len(all_ticks)}")
print(f"Output: {output_csv}")
if __name__ == "__main__":
batch_decode_instrument('./data/EURUSD', './EURUSD_ticks.csv')
Utilizzo:
python3 decode_bi5_batch.py
Errori comuni
- ⚠️ Il mese è indicizzato da zero nel percorso S3, ma non nei calcoli con le date — converti esplicitamente.
- ⚠️ Valore del punto errato: una discrepanza tra 100.000 e 1.000 corrompe silenziosamente i prezzi. Gli strumenti non-FX (indici/materie prime) non sono affatto coperti da questa regola — conferma esplicitamente il loro valore del punto (vedi
OTHER_POINT_VALUESsopra). - ⚠️ File mancanti ≠ errori: da trattare come "nessun tick quel giorno", non come un fallimento della decodifica.
- ⚠️ Flusso LZMA grezzo: alcune librerie richiedono una modalità "raw" esplicita — i decompressori
.xzstandard falliranno. - ⚠️ File orari obsoleti: i file più vecchi potrebbero usare millisecondi dall'inizio dell'ora invece di millisecondi dall'inizio del giorno. Nessuno dei decodificatori di questa guida rileva ciò automaticamente — entrambi presuppongono l'attuale formato giornaliero (
SYMBOL/YEAR/MONTH/DAY_ticks.bi5) in tutti i casi. Se stai lavorando con uno strumento o un intervallo di date sufficientemente vecchio da precedere la struttura giornaliera, verifica la base temporale effettiva del file prima di decodificare, invece di presumere che corrisponda al formato giornaliero qui descritto.
Pipeline end-to-end (download → decodifica → esportazione)
Combina download, decodifica ed esportazione in un unico script automatizzato — con filtraggio opzionale per data e pulizia.
Dipendenze
pip install boto3 pandas pyarrow
pyarrowè richiesto solo per l'esportazione in Parquet (consigliata per set di dati di grandi dimensioni).
Script completo della pipeline
#!/usr/bin/env python3
"""
Pipeline end-to-end: download -> decodifica -> esportazione
File tick giornalieri .bi5 in stile Dukascopy da 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}' non sembra una coppia FX standard a 6 lettere — "
f"verrà usato il valore del punto predefinito 100000. Verifica questo prima di fidarti del risultato."
)
return 100000
# Fase 1: Download
def download_instrument(instrument, destination, start_date=None, end_date=None):
"""Scarica tutti i file .bi5 di uno strumento, opzionalmente filtrati per intervallo di date."""
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"Elenco degli oggetti per {instrument}...")
objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
logger.info(f"Trovati {len(objects)} oggetti")
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"Impossibile analizzare la data dalla chiave: {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"✗ Download non riuscito per {obj.key}: {e}")
failed += 1
logger.info(f"Download completato: {downloaded} scaricati, {skipped} saltati (filtro data), {failed} falliti")
return downloaded, failed
# Fase 2: Decodifica
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):
"""Decodifica tutti i file .bi5 giornalieri di una cartella in un elenco di tuple di tick."""
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"✗ Errore durante la decodifica di {bi5_file}: {e}")
error_count += 1
all_ticks.sort(key=lambda t: t[0])
logger.info(f"Decodificati {file_count} file ({error_count} errori), {len(all_ticks)} tick totali")
return all_ticks
# Fase 3: Esportazione
def export_ticks(ticks, output_path, output_format='csv'):
"""Esporta i tick decodificati in CSV o Parquet utilizzando 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"✓ Esportati {len(df)} tick in {output_path} ({output_format})")
return df
# Orchestrazione della 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 avviata per {instrument} ===")
if start_date or end_date:
logger.info(f"Intervallo di date: {start_date or 'più remoto'} a {end_date or 'più recente'}")
downloaded, failed = download_instrument(instrument, raw_dir, start_date, end_date)
if failed:
logger.warning(f"{failed} file non scaricati per {instrument} — si procede con i {downloaded} scaricati con successo")
if downloaded == 0:
logger.warning(f"Nessun file scaricato per {instrument} — decodifica/esportazione saltata")
return
ticks = decode_instrument_folder(raw_dir, instrument)
if not ticks:
logger.warning(f"Nessun tick decodificato per {instrument}")
return
export_ticks(ticks, output_path, output_format)
if cleanup:
logger.info(f"Pulizia dei file .bi5 grezzi in {raw_dir}...")
shutil.rmtree(raw_dir, ignore_errors=True)
logger.info("✓ Pulizia completata")
logger.info(f"=== Pipeline completata per {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
)
Esempi di utilizzo
# Scaricare, decodificare ed esportare EURUSD in Parquet (predefinito)
python3 pipeline.py EURUSD
# Esportare invece in CSV
python3 pipeline.py EURUSD --format csv
# Filtrare per uno specifico intervallo di date
python3 pipeline.py EURUSD --start-date 2024-01-01 --end-date 2024-03-31
# Pulire i file .bi5 grezzi dopo l'esportazione
python3 pipeline.py EURUSD --cleanup
# Directory di output personalizzata
python3 pipeline.py GBPUSD --output-dir ./exports --cleanup
Pipeline in blocco per più strumenti
#!/bin/bash
# run_all_pipelines.sh
INSTRUMENTS=("EURUSD" "GBPUSD" "USDJPY" "AUDUSD")
for INSTRUMENT in "${INSTRUMENTS[@]}"; do
echo "=== Elaborazione di $INSTRUMENT ==="
python3 pipeline.py "$INSTRUMENT" --format parquet --cleanup
done
echo "✓ Tutte le pipeline sono state completate"
Best practice ✅
Download e controllo dei costi
- ✓ Includi sempre i flag
--region eu-west-1e--request-payer requester - ✓ Verifica la dimensione/numero di file effettivi per coppia con
--summarizeprima dei download in blocco — le medie possono essere fuorvianti (ad es. la dimensione media reale del file di EUR/USD è ~5 volte quella media dell'intero archivio) - ✓ Usa
aws s3 sync(noncp) per più file — vengono trasferiti solo i dati modificati - ✓ Usa il flag
--deleteper evitare di riscaricare file già esistenti - ✓ Esegui prima l'individuazione degli strumenti per confermare nomi di coppie esatti e validi — evita sincronizzazioni fallite a causa di errori di battitura
- ✓ Conserva i file di log per gli audit trail
- ✓ Configura allarmi CloudWatch per i download falliti
- ✓ Ricalcola i costi ogni volta che la dimensione del bucket/numero di file cambia in modo significativo
- ✓ Preferisci, quando possibile, il download di coppie specifiche invece dell'intero archivio
Decodifica e integrità dei dati
- ✓ Conferma sempre il valore del punto corretto (100.000 o 1.000) per ciascuno strumento prima della decodifica — una scalatura errata corrompe silenziosamente i prezzi
- ✓ Tratta i file
.bi5giornalieri mancanti come "nessun tick quel giorno" (fine settimana/festività), non come errori - ✓ Non mescolare mai vecchi file orari
.bi5con nuovi file giornalieri nella stessa pipeline senza prima rilevare il formato - ✓ Verifica la modalità di decompressione LZMA — alcune librerie richiedono una modalità "raw" esplicita (senza intestazioni del contenitore)
- ✓ Ordina cronologicamente i tick decodificati prima dell'esportazione, in caso di incongruenze nell'ordine del file system
Pipeline e automazione
- ✓ Usa
--cleanupper eliminare i file.bi5grezzi dopo la decodifica — possono sempre essere riscaricati da S3 se necessario - ✓ Preferisci Parquet a CSV per i dati a livello di tick — tipicamente 5-10 volte più piccolo e più veloce da interrogare
- ✓ Usa i filtri
--start-date/--end-datedurante i test, per evitare costi S3 non necessari - ✓ Combina con l'individuazione degli strumenti per scorrere automaticamente tutte le coppie valide invece di codificare elenchi fissi
- ✓ Registra ogni esecuzione della pipeline (conteggi scaricati/decodificati/esportati) per la verificabilità
Generale
- ✓ Esegui i download nelle ore di minore affluenza per ridurre la contesa (non il costo)
- ✓ Implementa checksum per la verifica dell'integrità dei dati dove è critico
- ✓ Monitora il volume totale di dati trasferiti nel tempo per controllare i costi cumulativi