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 / (o Delimiter='/' in Boto3) è essenziale — genera una risposta leggera CommonPrefixes invece 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 --delete per 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 / FileNotFoundError come "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_VALUES sopra).
  • ⚠️ 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 .xz standard 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-1 e --request-payer requester
  • ✓ Verifica la dimensione/numero di file effettivi per coppia con --summarize prima 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 (non cp) per più file — vengono trasferiti solo i dati modificati
  • ✓ Usa il flag --delete per 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 .bi5 giornalieri mancanti come "nessun tick quel giorno" (fine settimana/festività), non come errori
  • ✓ Non mescolare mai vecchi file orari .bi5 con 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 --cleanup per eliminare i file .bi5 grezzi 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-date durante 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

common.disclaimer