Datos históricos de precios
Resumen
Esta guía le permite acceder a los archivos de datos históricos de precios del bucket S3 cfg-public-proper-wallaby (eu-west-1), configurado con la política Requester Pays, garantizando un rendimiento óptimo de descarga y estabilidad del sistema.
Requisitos previos
- AWS CLI instalado y configurado
- Credenciales de AWS válidas con los permisos adecuados
- Acceso a la región eu-west-1
- Comprensión del modelo de facturación requester pays
Pasos de configuración
1. Establecer las credenciales de AWS y la región
export AWS_ACCESS_KEY_ID="your-access-key"
export AWS_SECRET_ACCESS_KEY="your-secret-key"
export AWS_DEFAULT_REGION="eu-west-1"
2. Verificar el acceso al bucket
aws s3 ls s3://cfg-public-proper-wallaby/ \
--region eu-west-1 \
--request-payer requester
3. Descubrir los instrumentos disponibles
Antes de descargar, liste todos los pares de divisas (instrumentos) disponibles en el bucket. Use --delimiter / (o el simple aws s3 ls) para listar solo los prefijos de nivel superior — evite el listado recursivo, que recorrería cada archivo solo para encontrar los ~20-30 nombres de carpetas.
Lista rápida (la más sencilla):
aws s3 ls s3://cfg-public-proper-wallaby/ \
--region eu-west-1 \
--request-payer requester
Lista limpia (solo nombres, ordenados, guardados en un archivo):
aws s3api list-objects-v2 \
--bucket cfg-public-proper-wallaby \
--delimiter "/" \
--region eu-west-1 \
--request-payer requester \
--query "CommonPrefixes[].Prefix" \
--output text | tr '\t' '\n' | sed 's|/$||' | sort > instruments.txt
echo "Total de instrumentos: $(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):
"""Obtiene todos los prefijos de instrumentos de nivel superior (pares de divisas)"""
instruments = []
paginator = s3_client.get_paginator('list_objects_v2')
page_iterator = paginator.paginate(
Bucket=bucket,
Delimiter='/',
RequestPayer='requester'
)
for page in page_iterator:
for prefix in page.get('CommonPrefixes', []):
instrument = prefix['Prefix'].rstrip('/')
instruments.append(instrument)
return sorted(instruments)
if __name__ == "__main__":
BUCKET = 'cfg-public-proper-wallaby'
instruments = list_instruments(BUCKET)
print(f"Total de instrumentos encontrados: {len(instruments)}\n")
for inst in instruments:
print(inst)
with open('instruments.txt', 'w') as f:
f.write('\n'.join(instruments))
Uso:
python3 list_instruments.py
⚠️ Nota sobre rendimiento y costo: El uso de
--delimiter /(oDelimiter='/'en Boto3) es esencial — genera una respuesta ligera deCommonPrefixesen lugar de enumerar cada archivo. Esto cuesta solo unas pocas solicitudes LIST económicas y se completa en segundos, frente a un escaneo completo del bucket, lento y más costoso, sin el delimitador.
4. Comando básico de acceso a S3 con Requester Pays
# Descargar un solo archivo
aws s3 cp s3://cfg-public-proper-wallaby/EURUSD/file.json . \
--region eu-west-1 \
--request-payer requester
5. Descarga por lotes con optimización de rendimiento
# Descargar todo el historial de un par de divisas específico
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
Parámetros de optimización de rendimiento
| Parámetro | Valor | Beneficio |
|---|---|---|
--max-concurrent-requests |
20-30 | Aumenta las descargas paralelas |
--max-bandwidth |
100MB/s | Evita los cuellos de botella de red |
--no-progress |
- | Reduce la sobrecarga de E/S |
--region |
eu-west-1 | Reduce la latencia dentro de la región de la UE |
6. Script de descarga avanzado (velocidad y estabilidad)
#!/bin/bash
BUCKET_NAME="cfg-public-proper-wallaby"
REGION="eu-west-1"
CURRENCY_PAIR="${1:-EURUSD}" # Predeterminado EURUSD, o pasado como argumento
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 "Iniciando sincronización S3 para $CURRENCY_PAIR desde $BUCKET_NAME en la región $REGION..." | tee "$LOG_FILE"
for attempt in $(seq 1 $MAX_RETRIES); do
echo "Intento de descarga $attempt de $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 "✓ Descarga completada con éxito el $(date)" | tee -a "$LOG_FILE"
echo "Total de archivos: $(find $DESTINATION -type f | wc -l)" | tee -a "$LOG_FILE"
echo "Tamaño total: $(du -sh $DESTINATION | cut -f1)" | tee -a "$LOG_FILE"
exit 0
else
echo "✗ Descarga fallida con código de salida $EXIT_CODE. Reintentando en 30 segundos..." | tee -a "$LOG_FILE"
sleep 30
fi
done
echo "✗ Descarga fallida después de $MAX_RETRIES intentos" | tee -a "$LOG_FILE"
exit 1
Uso:
chmod +x download-price-history.sh
./download-price-history.sh EURUSD
./download-price-history.sh GBPUSD
./download-price-history.sh
7. Descarga por lotes de varios pares de divisas
#!/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 "Iniciando descarga para $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 completado"
break
else
echo "⚠ $PAIR intento $attempt fallido, reintentando..."
sleep 30
fi
done
done
echo "✓ Todas las descargas completadas"
8. Uso de Python Boto3 (acceso programático)
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):
"""Descarga un archivo con Requester Pays habilitado"""
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"✓ Descargado: {key}")
return True
except Exception as e:
logger.error(f"✗ Error al descargar {key}: {e}")
return False
def batch_download(bucket, currency_pair, destination, max_workers=10):
"""Descarga todos los archivos de un par de divisas en paralelo"""
try:
bucket_obj = s3_resource.Bucket(bucket)
# Una barra diagonal final en el prefijo evita coincidir por error
# con otro instrumento que tenga este texto como prefijo
# (por ejemplo, "EURUSD" frente a un hipotético "EURUSDT").
prefix = f"{currency_pair}/"
objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
logger.info(f"Se encontraron {len(objects)} objetos para descargar de {currency_pair}")
with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
futures = []
for obj in objects:
if obj.key.endswith('/'): # Omitir directorios
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=== Resumen de descarga ===")
logger.info(f"Par de divisas: {currency_pair}")
logger.info(f"Total de archivos: {len(futures)}")
logger.info(f"Exitosos: {completed}")
logger.info(f"Fallidos: {failed}")
return completed, failed
except Exception as e:
logger.error(f"Error en la descarga por lotes: {e}")
return 0, len(objects)
# Uso - Par de divisas único
if __name__ == "__main__":
BUCKET = 'cfg-public-proper-wallaby'
CURRENCY_PAIR = 'EURUSD' # Modificar según sea necesario
DESTINATION = f'./data/{CURRENCY_PAIR}'
logger.info(f"Iniciando descarga desde s3://{BUCKET}/{CURRENCY_PAIR}/")
logger.info(f"Región: eu-west-1")
logger.info(f"Destino: {DESTINATION}")
completed, failed = batch_download(
bucket=BUCKET,
currency_pair=CURRENCY_PAIR,
destination=DESTINATION,
max_workers=10
)
if failed == 0:
logger.info("\n✓ ¡Todos los archivos se descargaron correctamente!")
else:
logger.warning(f"\n⚠ {failed} archivos no se pudieron descargar")
Uso:
python3 download_price_history.py
# Editar la variable CURRENCY_PAIR para otros pares
Nota: Dado que se confirmaron 26.586 archivos para EUR/USD, la llamada list(bucket_obj.objects.filter(...)) anterior los enumerará todos mediante llamadas paginadas a ListObjectsV2 antes de que comiencen las descargas. Este propio paso de listado genera su propio costo (pequeño) de solicitud, facturado por separado de las solicitudes GET. Para prefijos muy grandes, considere paginar manualmente o probar con recuperaciones parciales.
9. Descargar todos los pares de divisas (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):
"""Descarga un archivo con Requester Pays habilitado"""
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"✗ Error al descargar {key}: {e}")
return False
def download_all_pairs(bucket, destination_base, max_workers=10):
"""Descarga todo el historial de todos los pares de divisas"""
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"Se encontraron {len(pairs)} pares de divisas: {sorted(pairs)}")
total_completed = 0
total_failed = 0
for pair in sorted(pairs):
logger.info(f"\n=== Procesando {pair} ===")
# La barra diagonal final limita esto exactamente a las claves de este instrumento
prefix = f"{pair}/"
pair_objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
logger.info(f"Se encontraron {len(pair_objects)} objetos para {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)} archivos descargados")
logger.info(f"\n=== RESUMEN FINAL ===")
logger.info(f"Total completados: {total_completed}")
logger.info(f"Total fallidos: {total_failed}")
return total_completed, total_failed
except Exception as e:
logger.error(f"Error: {e}")
return 0, 0
# Uso
if __name__ == "__main__":
BUCKET = 'cfg-public-proper-wallaby'
DESTINATION_BASE = './data'
logger.info(f"Iniciando descarga del archivo completo desde s3://{BUCKET}/")
logger.info(f"Región: eu-west-1")
logger.info(f"Destino: {DESTINATION_BASE}")
download_all_pairs(
bucket=BUCKET,
destination_base=DESTINATION_BASE,
max_workers=10
)
10. Comandos de referencia rápida
# Listar todos los pares de divisas disponibles
aws s3 ls s3://cfg-public-proper-wallaby/ \
--region eu-west-1 \
--request-payer requester
# Descargar el historial de EURUSD
aws s3 sync s3://cfg-public-proper-wallaby/EURUSD/ ./EURUSD/ \
--region eu-west-1 \
--request-payer requester \
--max-concurrent-requests 25
# Descargar el historial de GBPUSD
aws s3 sync s3://cfg-public-proper-wallaby/GBPUSD/ ./GBPUSD/ \
--region eu-west-1 \
--request-payer requester \
--max-concurrent-requests 25
# Contar archivos en EURUSD
aws s3 ls s3://cfg-public-proper-wallaby/EURUSD/ \
--region eu-west-1 \
--request-payer requester \
--recursive | wc -l
# Obtener el tamaño total de EURUSD
aws s3 ls s3://cfg-public-proper-wallaby/EURUSD/ \
--region eu-west-1 \
--request-payer requester \
--recursive --summarize | grep "Total Size"
Puntos clave de motivación
Mejora de la velocidad ⚡
- Procesamiento paralelo: 20-25 conexiones simultáneas maximizan el rendimiento
- Optimización de EU-West-1: sin latencia entre regiones
- Estrategia de reintento adaptativa: recuperación inteligente ante fallos
- Agrupación de conexiones: minimiza la sobrecarga entre solicitudes
- Control de ancho de banda: evita la saturación de la red y problemas de estabilidad
Mejora de la estabilidad
- Mecanismo de reintento: 3 intentos con intervalos de 30 segundos
- Registro de errores: registros detallados para la resolución de problemas
- Tiempo de espera de conexión: 5 segundos para conectar, 60 segundos para lectura
- Modo adaptativo: ajusta la estrategia de reintento según el tipo de error
- Comprobaciones de integridad: validan las descargas de archivos
- Degradación controlada: continúa ante fallos parciales
Estimación de costos
Estructura de precios
Precios de solicitudes GET de AWS S3 (Requester Pays, eu-west-1):
- $0,0004 por cada 1.000 solicitudes
- $0,02 por GB de transferencia de datos (tarifa de salida de eu-west-1)
Los pasos de decodificación y de la pipeline en las secciones "Decodificación de archivos
.bi5" y "Pipeline de extremo a extremo" se ejecutan completamente en su equipo local — no generan costos adicionales de AWS más allá de los costos de descarga indicados aquí.
Estimación del archivo completo
⚠️ Archivo completo: ~400 GB, ~20.000.000 archivos
Costo de solicitudes = (20.000.000 / 1.000) × $0,0004 = $8,00
Costo de transferencia = 400 × $0,02 = $8,00
─────────────────────────────────────────────────────────
TOTAL = $16,00
Tamaño medio de archivo en todo el archivo: 400 GB / 20.000.000 archivos ≈ 20,97 KB/archivo
Ejemplo real: EUR/USD (datos reales verificados)
Confirmado mediante aws s3 ls --summarize:
Total de objetos: 26.586
Tamaño total: 2,6 GB
Tamaño medio de archivo: 2,6 GB / 26.586 ≈ 100,1 KB/archivo
Costos de solicitudes:
Número de solicitudes: 26.586
Costo por 1.000 solicitudes: $0,0004
Costo total de solicitudes = (26.586 / 1.000) × $0,0004
= 26,586 × $0,0004
≈ $0,0106
Costos de transferencia de datos:
Tamaño de los datos: 2,6 GB
Costo por GB: $0,02
Costo total de transferencia = 2,6 × $0,02 = $0,052
Costo total de descarga de EUR/USD (verificado):
Costos de solicitudes: $0,0106
Costos de transferencia de datos: $0,052
────────────────────────────
TOTAL: ~$0,06
Tabla comparativa de costos
| Escenario | Archivos | Tamaño | Costo de solicitudes | Costo de transferencia | Costo total |
|---|---|---|---|---|---|
| Archivo completo | 20M | 400 GB | $8,00 | $8,00 | $16,00 |
| EUR/USD (verificado, real) | 26.586 | 2,6 GB | $0,0106 | $0,052 | ~$0,06 |
| 5 pares (si similares a EURUSD) | ~133.000 | ~13 GB | $0,053 | $0,26 | ~$0,31 |
| 10 pares (si similares a EURUSD) | ~266.000 | ~26 GB | $0,106 | $0,52 | ~$0,63 |
Las estimaciones para "5 pares" y "10 pares" asumen un número de archivos/tamaño similar a EUR/USD. Los costos reales variarán según el par: los pares principales podrían tener más historial/archivos que EUR/USD, y los exóticos probablemente menos.
Consejos para optimizar costos
- ✓ Descargue solo los pares de divisas necesarios (no todo el archivo)
- ✓ Verifique siempre el tamaño/número de archivos reales por par con
--summarize— el tamaño medio real de archivo de EUR/USD (~100 KB) es aproximadamente 5 veces mayor que el promedio de todo el archivo (~21 KB), lo que demuestra lo engañosos que pueden ser los promedios - ✓ Agrupe las descargas para minimizar la sobrecarga de solicitudes
- ✓ Use el indicador
--deletepara evitar volver a descargar archivos ya existentes - ✓ A ~$0,06 por par (verificado con EUR/USD), descargar individualmente incluso decenas de pares sigue representando solo una pequeña fracción del costo del archivo completo de $16
- ✓ Use la Sección 3 (Descubrir los instrumentos disponibles) para confirmar los nombres exactos de los pares antes de ejecutar scripts por lotes
Acción recomendada antes de descargas masivas
aws s3 ls s3://cfg-public-proper-wallaby/<PAIR>/ \
--region eu-west-1 \
--request-payer requester \
--recursive --summarize | tail -3
Esto devuelve los valores exactos de Total Objects y Total Size, permitiendo calcular un costo preciso mediante:
Costo de solicitudes = (archivos / 1.000) × $0,0004
Costo de transferencia = tamaño_en_GB × $0,02
Decodificación de archivos .bi5 en datos tick legibles
Los archivos .bi5 de Dukascopy son archivos binarios de ticks comprimidos con LZMA. Con la estructura actual del bucket, cada archivo representa un día completo de datos tick por instrumento.
Estructura de archivos
Compresión: flujo LZMA sin procesar (no formato contenedor .xz)
Convención de rutas (diaria):
SYMBOL/YEAR/MONTH/DAY_ticks.bi5
Ejemplo: EURUSD/2024/00/15_ticks.bi5 → ticks de EUR/USD del 15 de enero de 2024
⚠️ El mes se indexa desde cero (enero =
00, diciembre =11).✅ Sin archivos vacíos: un archivo faltante para un día determinado significa que ese día no se registraron ticks (fines de semana, festivos). Trate las claves S3 faltantes /
FileNotFoundErrorcomo "sin datos", no como un error.
Formato del registro descomprimido
Cada tick es un registro binario fijo de 20 bytes, en formato big-endian:
| Bytes | Campo | Tipo | Notas |
|---|---|---|---|
| 0–3 | Marca de tiempo | uint32 | Milisegundos desde el inicio del día (UTC) |
| 4–7 | Precio Ask | uint32 | Requiere escalado por el valor de punto |
| 8–11 | Precio Bid | uint32 | Requiere escalado por el valor de punto |
| 12–15 | Volumen Ask | float32 | En millones de unidades de la divisa base |
| 16–19 | Volumen Bid | float32 | En millones de unidades de la divisa base |
Escalado de precios ("valor de punto")
| Tipo de instrumento | Valor de punto | Ejemplo |
|---|---|---|
| La mayoría de pares FX | 100.000 | EUR/USD, GBP/USD |
| Pares JPY | 1.000 | USD/JPY, EUR/JPY |
| Índices/materias primas | Varía | Verificar por instrumento |
⚠️ Las funciones auxiliares siguientes solo distinguen entre "par JPY" y "todo lo demás (100.000)". Para índices, materias primas o cualquier instrumento no FX, confirme el valor de punto correcto con el proveedor de datos antes de decodificar — aplicar silenciosamente 100.000 a un instrumento no FX producirá precios incorrectos sin generar ningún error. La función
get_point_value()actualizada en esta guía ahora emite una advertencia cuando recurre al valor predeterminado para un instrumento no reconocido y que no es JPY, para que esto no pase desapercibido.
Decodificador Python (día único)
import lzma
import struct
from datetime import datetime, timedelta
def decode_bi5_daily(filepath, day_start, point_value=100000):
"""
Decodifica un archivo tick .bi5 diario.
:param filepath: Ruta al archivo .bi5
:param day_start: datetime que representa las 00:00 UTC de ese día
:param point_value: 100000 para la mayoría de los pares, 1000 para los pares JPY
:return: Lista de diccionarios de ticks
"""
with open(filepath, 'rb') as f:
compressed = f.read()
raw = lzma.decompress(compressed)
ticks = []
record_size = 20
for i in range(0, len(raw), record_size):
chunk = raw[i:i+record_size]
ms, ask, bid, ask_vol, bid_vol = struct.unpack('>IIIff', chunk)
timestamp = day_start + timedelta(milliseconds=ms)
ticks.append({
'time': timestamp,
'ask': ask / point_value,
'bid': bid / point_value,
'ask_volume': ask_vol,
'bid_volume': bid_vol
})
return ticks
# Ejemplo de uso
day_start = datetime(2024, 1, 15, 0, 0, 0)
ticks = decode_bi5_daily('EURUSD/2024/00/15_ticks.bi5', day_start, point_value=100000)
print(f"Total de ticks del día: {len(ticks)}")
for tick in ticks[:5]:
print(tick)
Script de decodificación por lotes (carpeta de instrumento completa → 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'}
# Instrumentos no FX conocidos (índices/materias primas) con un valor de punto confirmado.
# Amplíe esta lista a medida que confirme otros instrumentos con el proveedor de datos.
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}' no parece un par FX estándar de 6 letras y no tiene "
f"un valor de punto confirmado — se usará el valor predeterminado 100000. Verifique esto antes de confiar en el resultado."
)
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 todos los archivos .bi5 diarios de un instrumento en un único CSV.
Espera la estructura de carpetas: 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"✗ Error al decodificar {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=== Resumen de decodificación ===")
print(f"Instrumento: {instrument_name}")
print(f"Archivos procesados: {file_count}")
print(f"Errores: {error_count}")
print(f"Total de ticks: {len(all_ticks)}")
print(f"Salida: {output_csv}")
if __name__ == "__main__":
batch_decode_instrument('./data/EURUSD', './EURUSD_ticks.csv')
Uso:
python3 decode_bi5_batch.py
Errores comunes
- ⚠️ El mes se indexa desde cero en la ruta S3, pero no en los cálculos de fechas — conviértalo explícitamente.
- ⚠️ Valor de punto incorrecto: una discrepancia entre 100.000 y 1.000 corrompe silenciosamente los precios. Los instrumentos no FX (índices/materias primas) no están cubiertos en absoluto por esta regla — confirme su valor de punto explícitamente (véase
OTHER_POINT_VALUESarriba). - ⚠️ Archivos faltantes ≠ errores: trátelos como "sin ticks ese día", no como un fallo de decodificación.
- ⚠️ Flujo LZMA sin procesar: algunas bibliotecas requieren un modo "raw" explícito — los descompresores
.xzestándar fallarán. - ⚠️ Archivos horarios antiguos: los archivos más antiguos podrían usar milisegundos desde el inicio de la hora en lugar de milisegundos desde el inicio del día. Ninguno de los decodificadores de esta guía detecta esto automáticamente — ambos asumen el formato diario actual (
SYMBOL/YEAR/MONTH/DAY_ticks.bi5) en todo momento. Si trabaja con un instrumento o rango de fechas lo suficientemente antiguo como para preceder al diseño diario, verifique la base temporal real del archivo antes de decodificar, en lugar de asumir que coincide con el formato diario descrito aquí.
Pipeline de extremo a extremo (descarga → decodificación → exportación)
Combina la descarga, la decodificación y la exportación en un único script automatizado — con filtrado opcional por fecha y limpieza.
Dependencias
pip install boto3 pandas pyarrow
pyarrowsolo es necesario para la exportación a Parquet (recomendado para conjuntos de datos grandes).
Script completo de la pipeline
#!/usr/bin/env python3
"""
Pipeline de extremo a extremo: descarga -> decodificación -> exportación
Archivos tick diarios .bi5 estilo Dukascopy desde 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}' no parece un par FX estándar de 6 letras — "
f"se usará el valor de punto predeterminado 100000. Verifique esto antes de confiar en el resultado."
)
return 100000
# Paso 1: Descarga
def download_instrument(instrument, destination, start_date=None, end_date=None):
"""Descarga todos los archivos .bi5 de un instrumento, opcionalmente filtrados por rango de fechas."""
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"Listando objetos para {instrument}...")
objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
logger.info(f"Se encontraron {len(objects)} objetos")
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"No se pudo analizar la fecha de la clave: {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"✗ Error al descargar {obj.key}: {e}")
failed += 1
logger.info(f"Descarga completada: {downloaded} descargados, {skipped} omitidos (filtro de fecha), {failed} fallidos")
return downloaded, failed
# Paso 2: Decodificación
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 todos los archivos .bi5 diarios de una carpeta en una lista de tuplas de ticks."""
instrument_dir = Path(instrument_dir)
point_value = get_point_value(instrument_name)
all_ticks = []
file_count = 0
error_count = 0
for year_dir in sorted(instrument_dir.iterdir()):
if not year_dir.is_dir():
continue
for month_dir in sorted(year_dir.iterdir()):
if not month_dir.is_dir():
continue
year = int(year_dir.name)
month = int(month_dir.name) + 1
for bi5_file in sorted(month_dir.glob('*_ticks.bi5')):
try:
day = int(bi5_file.name.split('_')[0])
day_start = datetime(year, month, day)
ticks = decode_bi5_daily(bi5_file, day_start, point_value)
all_ticks.extend(ticks)
file_count += 1
except Exception as e:
logger.error(f"✗ Error al decodificar {bi5_file}: {e}")
error_count += 1
all_ticks.sort(key=lambda t: t[0])
logger.info(f"Decodificados {file_count} archivos ({error_count} errores), {len(all_ticks)} ticks en total")
return all_ticks
# Paso 3: Exportación
def export_ticks(ticks, output_path, output_format='csv'):
"""Exporta los ticks decodificados a CSV o Parquet usando 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"✓ Se exportaron {len(df)} ticks a {output_path} ({output_format})")
return df
# Orquestación de la 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 iniciada para {instrument} ===")
if start_date or end_date:
logger.info(f"Rango de fechas: {start_date or 'la más temprana'} a {end_date or 'la más reciente'}")
downloaded, failed = download_instrument(instrument, raw_dir, start_date, end_date)
if failed:
logger.warning(f"{failed} archivo(s) no se pudieron descargar para {instrument} — se continuará con los {downloaded} descargados con éxito")
if downloaded == 0:
logger.warning(f"No se descargó ningún archivo para {instrument} — se omite la decodificación/exportación")
return
ticks = decode_instrument_folder(raw_dir, instrument)
if not ticks:
logger.warning(f"No se decodificó ningún tick para {instrument}")
return
export_ticks(ticks, output_path, output_format)
if cleanup:
logger.info(f"Limpiando los archivos .bi5 originales en {raw_dir}...")
shutil.rmtree(raw_dir, ignore_errors=True)
logger.info("✓ Limpieza completada")
logger.info(f"=== Pipeline finalizada para {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
)
Ejemplos de uso
# Descargar, decodificar y exportar EURUSD a Parquet (predeterminado)
python3 pipeline.py EURUSD
# Exportar a CSV en su lugar
python3 pipeline.py EURUSD --format csv
# Filtrar por un rango de fechas específico
python3 pipeline.py EURUSD --start-date 2024-01-01 --end-date 2024-03-31
# Limpiar los archivos .bi5 originales después de exportar
python3 pipeline.py EURUSD --cleanup
# Directorio de salida personalizado
python3 pipeline.py GBPUSD --output-dir ./exports --cleanup
Pipeline por lotes para múltiples instrumentos
#!/bin/bash
# run_all_pipelines.sh
INSTRUMENTS=("EURUSD" "GBPUSD" "USDJPY" "AUDUSD")
for INSTRUMENT in "${INSTRUMENTS[@]}"; do
echo "=== Procesando $INSTRUMENT ==="
python3 pipeline.py "$INSTRUMENT" --format parquet --cleanup
done
echo "✓ Todas las pipelines completadas"
Buenas prácticas ✅
Descarga y control de costos
- ✓ Incluya siempre los indicadores
--region eu-west-1y--request-payer requester - ✓ Verifique el tamaño/número de archivos reales por par con
--summarizeantes de descargas masivas — los promedios pueden ser engañosos (por ejemplo, el tamaño medio real de archivo de EUR/USD es ~5 veces el promedio de todo el archivo) - ✓ Use
aws s3 sync(nocp) para varios archivos — solo se transfieren los datos modificados - ✓ Use el indicador
--deletepara evitar volver a descargar archivos ya existentes - ✓ Ejecute primero el descubrimiento de instrumentos para confirmar nombres de pares exactos y válidos — evita sincronizaciones fallidas por errores tipográficos
- ✓ Conserve los archivos de registro para auditorías
- ✓ Configure alarmas de CloudWatch para las descargas fallidas
- ✓ Recalcule los costos cada vez que el tamaño del bucket/número de archivos cambie de forma significativa
- ✓ Cuando sea posible, prefiera descargar pares específicos en lugar del archivo completo
Decodificación e integridad de datos
- ✓ Confirme siempre el valor de punto correcto (100.000 o 1.000) por instrumento antes de decodificar — un escalado incorrecto corrompe silenciosamente los precios
- ✓ Trate los archivos
.bi5diarios faltantes como "sin ticks ese día" (fines de semana/festivos), no como errores - ✓ Nunca mezcle archivos horarios antiguos
.bi5con archivos diarios nuevos en la misma pipeline sin detectar primero el formato - ✓ Verifique el modo de descompresión LZMA — algunas bibliotecas requieren un modo "raw" explícito (sin encabezados de contenedor)
- ✓ Ordene cronológicamente los ticks decodificados antes de exportar, en caso de inconsistencias en el orden del sistema de archivos
Pipeline y automatización
- ✓ Use
--cleanuppara eliminar los archivos.bi5originales después de decodificar — siempre se pueden volver a descargar de S3 si es necesario - ✓ Prefiera Parquet sobre CSV para datos a nivel de tick — normalmente 5-10 veces más pequeño y más rápido de consultar
- ✓ Use los filtros
--start-date/--end-datedurante las pruebas, para evitar costos innecesarios de S3 - ✓ Combine con el descubrimiento de instrumentos para recorrer automáticamente todos los pares válidos en lugar de codificar listas fijas
- ✓ Registre cada ejecución de la pipeline (recuentos de descargados/decodificados/exportados) para fines de auditoría
General
- ✓ Ejecute las descargas en horas de menor actividad para reducir la contención (no el costo)
- ✓ Implemente sumas de verificación para comprobar la integridad de los datos donde sea crítico
- ✓ Monitoree el volumen total de datos transferidos a lo largo del tiempo para controlar los costos acumulados