历史价格数据

概述

本指南将帮助您从配置了 Requester Pays(请求者付费)策略的 cfg-public-proper-wallaby S3 存储桶(eu-west-1 区域)中访问历史价格数据文件,确保最佳的下载性能和系统稳定性。

前提条件

  • 已安装并配置 AWS CLI
  • 拥有相应权限的有效 AWS 凭证
  • 可访问 eu-west-1 区域
  • 了解 requester pays 计费模式

配置步骤

1. 设置 AWS 凭证和区域

export AWS_ACCESS_KEY_ID="your-access-key"
export AWS_SECRET_ACCESS_KEY="your-secret-key"
export AWS_DEFAULT_REGION="eu-west-1"

2. 验证存储桶访问权限

aws s3 ls s3://cfg-public-proper-wallaby/ \
  --region eu-west-1 \
  --request-payer requester

3. 发现可用的交易品种

在下载之前,请先列出存储桶中所有可用的货币对(交易品种)。使用 --delimiter /(或直接使用 aws s3 ls)来仅列出顶层前缀 —— 请避免使用递归列表,因为那会扫描每一个文件,仅仅是为了找到约 20-30 个文件夹名称。

快速列表(最简单):

aws s3 ls s3://cfg-public-proper-wallaby/ \
  --region eu-west-1 \
  --request-payer requester

清晰列表(仅名称,已排序,保存至文件):

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 "交易品种总数: $(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):
    """获取所有顶层交易品种前缀(货币对)"""
    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"共找到交易品种: {len(instruments)}\n")
    for inst in instruments:
        print(inst)

    with open('instruments.txt', 'w') as f:
        f.write('\n'.join(instruments))

使用方法:

python3 list_instruments.py

⚠️ 性能与成本说明: 使用 --delimiter /(或 Boto3 中的 Delimiter='/')至关重要 —— 它会触发轻量级的 CommonPrefixes 响应,而不是枚举每一个文件。这只需花费少量低成本的 LIST 请求,并在几秒钟内完成,而不使用分隔符进行全存储桶扫描则速度慢且成本更高。

4. 使用 Requester Pays 的基本 S3 访问命令

# 下载单个文件
aws s3 cp s3://cfg-public-proper-wallaby/EURUSD/file.json . \
  --region eu-west-1 \
  --request-payer requester

5. 带性能优化的批量下载

# 下载特定货币对的全部历史数据
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

性能优化参数

参数 优势
--max-concurrent-requests 20-30 增加并行下载数
--max-bandwidth 100MB/s 防止网络瓶颈
--no-progress - 降低 I/O 开销
--region eu-west-1 降低欧盟区域内的延迟

6. 高级下载脚本(速度与稳定性)

#!/bin/bash

BUCKET_NAME="cfg-public-proper-wallaby"
REGION="eu-west-1"
CURRENCY_PAIR="${1:-EURUSD}"  # 默认为 EURUSD,也可作为参数传入
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 "正在从 $BUCKET_NAME(区域:$REGION)开始同步 $CURRENCY_PAIR 的 S3 数据..." | tee "$LOG_FILE"

for attempt in $(seq 1 $MAX_RETRIES); do
    echo "第 $attempt 次下载尝试,共 $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 "✓ 下载于 $(date) 成功完成" | tee -a "$LOG_FILE"
        echo "文件总数: $(find $DESTINATION -type f | wc -l)" | tee -a "$LOG_FILE"
        echo "总大小: $(du -sh $DESTINATION | cut -f1)" | tee -a "$LOG_FILE"
        exit 0
    else
        echo "✗ 下载失败,退出代码 $EXIT_CODE。30 秒后重试..." | tee -a "$LOG_FILE"
        sleep 30
    fi
done

echo "✗ 经过 $MAX_RETRIES 次尝试后下载失败" | tee -a "$LOG_FILE"
exit 1

使用方法:

chmod +x download-price-history.sh

./download-price-history.sh EURUSD
./download-price-history.sh GBPUSD
./download-price-history.sh

7. 批量下载多个货币对

#!/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 "开始下载 $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 已完成"
            break
        else
            echo "⚠ $PAIR 第 $attempt 次尝试失败,正在重试..."
            sleep 30
        fi
    done
done

echo "✓ 全部下载已完成"

8. 使用 Python Boto3(编程方式访问)

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):
    """在启用 Requester Pays 的情况下下载文件"""
    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"✓ 已下载: {key}")
        return True
    except Exception as e:
        logger.error(f"✗ 下载 {key} 时出错: {e}")
        return False

def batch_download(bucket, currency_pair, destination, max_workers=10):
    """并行下载某个货币对的所有文件"""
    try:
        bucket_obj = s3_resource.Bucket(bucket)

        # 前缀末尾加斜杠可避免意外匹配到恰好以此字符串为前缀的
        # 其他交易品种(例如 "EURUSD" 与假设存在的 "EURUSDT")。
        prefix = f"{currency_pair}/"
        objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
        logger.info(f"找到 {len(objects)} 个待下载的 {currency_pair} 对象")

        with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
            futures = []

            for obj in objects:
                if obj.key.endswith('/'):  # 跳过目录
                    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=== 下载摘要 ===")
            logger.info(f"货币对: {currency_pair}")
            logger.info(f"文件总数: {len(futures)}")
            logger.info(f"成功: {completed}")
            logger.info(f"失败: {failed}")

            return completed, failed

    except Exception as e:
        logger.error(f"批量下载错误: {e}")
        return 0, len(objects)

# 使用方法 - 单个货币对
if __name__ == "__main__":
    BUCKET = 'cfg-public-proper-wallaby'
    CURRENCY_PAIR = 'EURUSD'  # 根据需要修改
    DESTINATION = f'./data/{CURRENCY_PAIR}'

    logger.info(f"正在从 s3://{BUCKET}/{CURRENCY_PAIR}/ 开始下载")
    logger.info(f"区域: eu-west-1")
    logger.info(f"目标位置: {DESTINATION}")

    completed, failed = batch_download(
        bucket=BUCKET,
        currency_pair=CURRENCY_PAIR,
        destination=DESTINATION,
        max_workers=10
    )

    if failed == 0:
        logger.info("\n✓ 所有文件均已成功下载!")
    else:
        logger.warning(f"\n⚠ {failed} 个文件下载失败")

使用方法:

python3 download_price_history.py
# 编辑 CURRENCY_PAIR 变量以下载其他货币对

说明: 鉴于已确认 EUR/USD 共有 26,586 个文件,上面的 list(bucket_obj.objects.filter(...)) 调用将在任何下载开始之前,通过分页的 ListObjectsV2 调用枚举所有这些文件。此列表步骤本身会产生独立的(少量)请求成本,与 GET 请求分开计费。对于非常大的前缀,请考虑手动分页或先使用部分数据进行测试。

9. 下载所有货币对(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):
    """在启用 Requester Pays 的情况下下载文件"""
    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"✗ 下载 {key} 时出错: {e}")
        return False

def download_all_pairs(bucket, destination_base, max_workers=10):
    """下载所有货币对的全部历史数据"""
    try:
        bucket_obj = s3_resource.Bucket(bucket)

        pairs = set()
        for obj in bucket_obj.objects.filter(RequestPayer='requester'):
            pair = obj.key.split('/')[0]
            if pair and not obj.key.endswith('/'):
                pairs.add(pair)

        logger.info(f"找到 {len(pairs)} 个货币对: {sorted(pairs)}")

        total_completed = 0
        total_failed = 0

        for pair in sorted(pairs):
            logger.info(f"\n=== 正在处理 {pair} ===")

            # 前缀末尾的斜杠将范围精确限定为该交易品种的键
            prefix = f"{pair}/"
            pair_objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
            logger.info(f"找到 {len(pair_objects)} 个 {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)} 个文件")

        logger.info(f"\n=== 最终摘要 ===")
        logger.info(f"总计完成: {total_completed}")
        logger.info(f"总计失败: {total_failed}")

        return total_completed, total_failed

    except Exception as e:
        logger.error(f"错误: {e}")
        return 0, 0

# 使用方法
if __name__ == "__main__":
    BUCKET = 'cfg-public-proper-wallaby'
    DESTINATION_BASE = './data'

    logger.info(f"正在从 s3://{BUCKET}/ 开始下载完整存档")
    logger.info(f"区域: eu-west-1")
    logger.info(f"目标位置: {DESTINATION_BASE}")

    download_all_pairs(
        bucket=BUCKET,
        destination_base=DESTINATION_BASE,
        max_workers=10
    )

10. 快速参考命令

# 列出所有可用的货币对
aws s3 ls s3://cfg-public-proper-wallaby/ \
  --region eu-west-1 \
  --request-payer requester

# 下载 EURUSD 历史数据
aws s3 sync s3://cfg-public-proper-wallaby/EURUSD/ ./EURUSD/ \
  --region eu-west-1 \
  --request-payer requester \
  --max-concurrent-requests 25

# 下载 GBPUSD 历史数据
aws s3 sync s3://cfg-public-proper-wallaby/GBPUSD/ ./GBPUSD/ \
  --region eu-west-1 \
  --request-payer requester \
  --max-concurrent-requests 25

# 统计 EURUSD 中的文件数
aws s3 ls s3://cfg-public-proper-wallaby/EURUSD/ \
  --region eu-west-1 \
  --request-payer requester \
  --recursive | wc -l

# 获取 EURUSD 的总大小
aws s3 ls s3://cfg-public-proper-wallaby/EURUSD/ \
  --region eu-west-1 \
  --request-payer requester \
  --recursive --summarize | grep "Total Size"

关键要点

速度提升 ⚡

  • 并行处理:20-25 个并发连接可最大化吞吐量
  • EU-West-1 优化:无跨区域延迟
  • 自适应重试策略:智能故障恢复
  • 连接池:最大限度地减少请求之间的开销
  • 带宽控制:防止网络饱和和稳定性问题

稳定性提升

  • 重试机制:3 次尝试,间隔 30 秒
  • 错误日志记录:详细日志便于排查问题
  • 连接超时:连接超时 5 秒,读取超时 60 秒
  • 自适应模式:根据错误类型调整重试策略
  • 完整性检查:验证文件下载
  • 优雅降级:在部分失败时继续运行

成本估算

定价结构

AWS S3 GET 请求定价(Requester Pays,eu-west-1):

  • 每 1,000 次请求 $0.0004
  • 每 GB 数据传输 $0.02(eu-west-1 出口费率)

"解码 .bi5 文件"和"完整流水线"两节中的解码与流水线步骤完全在您的本地计算机上运行 —— 除了此处列出的下载成本外,它们不会产生任何额外的 AWS 费用。

完整存档估算

⚠️ 完整存档:约 400 GB,约 2000 万个文件

请求成本  = (20,000,000 / 1,000) × $0.0004 = $8.00
传输成本 = 400 × $0.02                      = $8.00
─────────────────────────────────────────────────────────
总计                                          = $16.00

整个存档的平均文件大小:400 GB / 2000 万个文件 ≈ 20.97 KB/文件

实际案例:EUR/USD(已验证的实际数据)

通过 aws s3 ls --summarize 确认:

对象总数: 26,586
总大小:    2.6 GB
平均文件大小: 2.6 GB / 26,586 ≈ 100.1 KB/文件

请求成本:

请求数量: 26,586
每 1,000 次请求成本: $0.0004

总请求成本 = (26,586 / 1,000) × $0.0004
                   = 26.586 × $0.0004
                   ≈ $0.0106

数据传输成本:

数据大小: 2.6 GB
每 GB 成本: $0.02
总传输成本 = 2.6 × $0.02 = $0.052

EUR/USD 下载总成本(已验证):

请求成本:        $0.0106
数据传输成本:  $0.052
────────────────────────────
总计:                ~$0.06

成本对比表

场景 文件数 大小 请求成本 传输成本 总成本
完整存档 2000万 400 GB $8.00 $8.00 $16.00
EUR/USD(已验证的实际数据) 26,586 2.6 GB $0.0106 $0.052 ~$0.06
5 个货币对(若与 EURUSD 相似) ~133,000 ~13 GB $0.053 $0.26 ~$0.31
10 个货币对(若与 EURUSD 相似) ~266,000 ~26 GB $0.106 $0.52 ~$0.63

"5 个货币对"和"10 个货币对"的估算假设文件数量/大小与 EUR/USD 相似。实际成本会因货币对而异 —— 主要货币对的历史数据/文件数可能比 EUR/USD 更多,而冷门货币对可能更少。

成本优化建议

  • ✓ 仅下载所需的货币对(而非整个存档)
  • 始终使用 --summarize 验证每个货币对的实际大小/文件数量 —— EUR/USD 的实际平均文件大小(约 100 KB)约为整个存档平均值(约 21 KB)的 5 倍,这表明平均值可能具有误导性
  • ✓ 批量处理下载以最大限度地减少请求开销
  • ✓ 使用 --delete 标志以避免重新下载已存在的文件
  • ✓ 按每个货币对约 $0.06 的成本(以 EUR/USD 验证为准)计算,即使单独下载数十个货币对,也仅占整个存档 $16 成本的一小部分
  • ✓ 在运行批量脚本之前,使用第 3 节(发现可用交易品种)确认准确的货币对名称

批量下载前的建议操作

aws s3 ls s3://cfg-public-proper-wallaby/<PAIR>/ \
  --region eu-west-1 \
  --request-payer requester \
  --recursive --summarize | tail -3

这将返回准确的 Total ObjectsTotal Size 值,从而可以通过以下公式计算准确的成本:

请求成本  = (文件数 / 1000) × $0.0004
传输成本 = 大小_GB × $0.02

解码 .bi5 文件为可读的逐笔数据

Dukascopy 的 .bi5 文件是经过 LZMA 压缩的二进制逐笔数据文件。按照当前的存储桶结构,每个文件代表每个交易品种一整天的逐笔数据。

文件结构

压缩方式: 原始 LZMA 流(不是 .xz 容器格式)

路径约定(按天):

SYMBOL/YEAR/MONTH/DAY_ticks.bi5

示例:EURUSD/2024/00/15_ticks.bi5 → 2024 年 1 月 15 日的 EUR/USD 逐笔数据

⚠️ 月份从零开始编号(一月 = 00,十二月 = 11)。

没有空文件:某一天没有对应文件,意味着当天没有记录任何成交(周末、节假日)。请将缺失的 S3 键 / FileNotFoundError 视为"无数据",而不是错误。

解压后的记录格式

每笔逐笔数据是固定 20 字节的二进制记录,采用大端字节序:

字节 字段 类型 备注
0–3 时间戳 uint32 自当天开始的毫秒数(UTC)
4–7 卖价(Ask) uint32 需要按点值缩放
8–11 买价(Bid) uint32 需要按点值缩放
12–15 卖出量(Ask) float32 以百万计的基础货币单位
16–19 买入量(Bid) float32 以百万计的基础货币单位

价格缩放("点值")

交易品种类型 点值 示例
大多数外汇对 100,000 EUR/USD, GBP/USD
JPY 相关货币对 1,000 USD/JPY, EUR/JPY
指数/大宗商品 各不相同 需按具体品种核实

⚠️ 下方的辅助函数仅区分"JPY 货币对"和"其他所有品种(100,000)"两种情况。对于指数、大宗商品或任何非外汇交易品种,请在解码前向数据提供方核实正确的点值 —— 若默默地对非外汇品种应用 100,000,会在不产生任何错误提示的情况下生成错误的价格。本指南更新后的 get_point_value() 函数现在会在对一个无法识别、且非 JPY 相关的品种回退使用默认值时输出警告,以避免这种情况被悄然忽略。

Python 解码器(单日)

import lzma
import struct
from datetime import datetime, timedelta

def decode_bi5_daily(filepath, day_start, point_value=100000):
    """
    解码单日 .bi5 逐笔数据文件。

    :param filepath: .bi5 文件的路径
    :param day_start: 表示该日 UTC 00:00 的 datetime
    :param point_value: 大多数货币对为 100000,JPY 相关货币对为 1000
    :return: 逐笔数据字典列表
    """
    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

# 使用示例
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"当日逐笔数据总数: {len(ticks)}")
for tick in ticks[:5]:
    print(tick)

批量解码脚本(完整交易品种文件夹 → 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'}

# 已确认点值的非外汇品种(指数/大宗商品)。
# 在向数据提供方核实其他品种后,请扩展此列表。
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}' 看起来不像标准的 6 字母外汇货币对,且没有已确认的点值 —— "
            f"将使用默认值 100000。请在信任结果前核实此项。"
        )
    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):
    """
    将某交易品种的所有每日 .bi5 文件解码为单个 CSV 文件。
    期望的文件夹结构为: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"✗ 解码 {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=== 解码摘要 ===")
    print(f"交易品种: {instrument_name}")
    print(f"已处理文件数: {file_count}")
    print(f"错误数: {error_count}")
    print(f"逐笔数据总数: {len(all_ticks)}")
    print(f"输出文件: {output_csv}")

if __name__ == "__main__":
    batch_decode_instrument('./data/EURUSD', './EURUSD_ticks.csv')

使用方法:

python3 decode_bi5_batch.py

常见陷阱

  • ⚠️ 月份在 S3 路径中从零开始编号,但在日期计算中并非如此 —— 需要显式转换。
  • ⚠️ 点值错误:100,000 与 1,000 之间的混淆会在无声无息中破坏价格数据。非外汇品种(指数/大宗商品)完全不受此规则覆盖 —— 请显式确认它们的点值(参见上文 OTHER_POINT_VALUES)。
  • ⚠️ 缺失文件 ≠ 错误:应将其视为"当天没有逐笔数据",而非解码失败。
  • ⚠️ 原始 LZMA 流:某些库需要显式的"raw"模式 —— 标准的 .xz 解压工具会失败。
  • ⚠️ 旧版按小时文件:较旧的文件可能使用自当小时开始的毫秒数而非自当天开始的毫秒数。本指南中的两种解码器都不会自动检测这一点 —— 两者始终假定采用当前的每日格式(SYMBOL/YEAR/MONTH/DAY_ticks.bi5)。如果您处理的交易品种或日期范围早于按日布局出现的时间,请在解码前核实文件的实际时间基准,而不要想当然地假定它符合本文所述的每日格式。

完整流水线(下载 → 解码 → 导出)

将下载、解码和导出整合到一个自动化脚本中 —— 支持可选的日期过滤和清理功能。

依赖项

pip install boto3 pandas pyarrow

仅在导出为 Parquet 格式时才需要 pyarrow(建议用于大型数据集)。

完整流水线脚本

#!/usr/bin/env python3
"""
端到端流水线:下载 -> 解码 -> 导出
来自 S3(Requester Pays)的 Dukascopy 风格每日 .bi5 逐笔数据文件
"""

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}' 看起来不像标准的 6 字母外汇货币对 —— "
            f"将使用默认点值 100000。请在信任结果前核实此项。"
        )
    return 100000

# 第 1 步:下载

def download_instrument(instrument, destination, start_date=None, end_date=None):
    """下载某交易品种的所有 .bi5 文件,可选按日期范围过滤。"""
    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"正在列出 {instrument} 的对象...")
    objects = list(bucket_obj.objects.filter(Prefix=prefix, RequestPayer='requester'))
    logger.info(f"找到 {len(objects)} 个对象")

    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"无法从键中解析日期: {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"✗ 下载 {obj.key} 失败: {e}")
            failed += 1

    logger.info(f"下载完成:已下载 {downloaded} 个,因日期过滤跳过 {skipped} 个,失败 {failed} 个")
    return downloaded, failed

# 第 2 步:解码

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):
    """将某文件夹中的所有每日 .bi5 文件解码为逐笔数据元组列表。"""
    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"✗ 解码 {bi5_file} 时出错: {e}")
                    error_count += 1

    all_ticks.sort(key=lambda t: t[0])
    logger.info(f"已解码 {file_count} 个文件({error_count} 个错误),共 {len(all_ticks)} 条逐笔数据")

    return all_ticks

# 第 3 步:导出

def export_ticks(ticks, output_path, output_format='csv'):
    """使用 pandas 将解码后的逐笔数据导出为 CSV 或 Parquet。"""
    df = pd.DataFrame(ticks, columns=['timestamp', 'ask', 'bid', 'ask_volume', 'bid_volume'])

    output_path = Path(output_path)
    output_path.parent.mkdir(parents=True, exist_ok=True)

    if output_format == 'csv':
        df.to_csv(output_path, index=False)
    elif output_format == 'parquet':
        df.to_parquet(output_path, index=False, engine='pyarrow', compression='snappy')
    else:
        raise ValueError(f"Unsupported format: {output_format}")

    logger.info(f"✓ 已将 {len(df)} 条逐笔数据导出至 {output_path}({output_format})")
    return df

# 流水线编排

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"=== {instrument} 的流水线已启动 ===")
    if start_date or end_date:
        logger.info(f"日期范围:{start_date or '最早'} 至 {end_date or '最晚'}")

    downloaded, failed = download_instrument(instrument, raw_dir, start_date, end_date)
    if failed:
        logger.warning(f"{instrument} 有 {failed} 个文件下载失败 —— 将继续处理已成功的 {downloaded} 个文件")
    if downloaded == 0:
        logger.warning(f"未为 {instrument} 下载任何文件 —— 跳过解码/导出")
        return

    ticks = decode_instrument_folder(raw_dir, instrument)
    if not ticks:
        logger.warning(f"未能为 {instrument} 解码出任何逐笔数据")
        return

    export_ticks(ticks, output_path, output_format)

    if cleanup:
        logger.info(f"正在清理 {raw_dir} 下的原始 .bi5 文件...")
        shutil.rmtree(raw_dir, ignore_errors=True)
        logger.info("✓ 清理完成")

    logger.info(f"=== {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
    )

使用示例

# 下载、解码并导出 EURUSD 为 Parquet(默认格式)
python3 pipeline.py EURUSD

# 改为导出为 CSV
python3 pipeline.py EURUSD --format csv

# 按指定日期范围过滤
python3 pipeline.py EURUSD --start-date 2024-01-01 --end-date 2024-03-31

# 导出后清理原始 .bi5 文件
python3 pipeline.py EURUSD --cleanup

# 自定义输出目录
python3 pipeline.py GBPUSD --output-dir ./exports --cleanup

多交易品种的批量流水线

#!/bin/bash
# run_all_pipelines.sh

INSTRUMENTS=("EURUSD" "GBPUSD" "USDJPY" "AUDUSD")

for INSTRUMENT in "${INSTRUMENTS[@]}"; do
    echo "=== 正在处理 $INSTRUMENT ==="
    python3 pipeline.py "$INSTRUMENT" --format parquet --cleanup
done

echo "✓ 所有流水线已完成"

最佳实践 ✅

下载与成本控制

  • ✓ 始终包含 --region eu-west-1--request-payer requester 参数
  • ✓ 在批量下载前,使用 --summarize 验证每个货币对的实际大小/文件数量 —— 平均值可能具有误导性(例如,EUR/USD 的实际平均文件大小约为整个存档平均值的 5 倍)
  • ✓ 对多个文件使用 aws s3 sync(而非 cp)—— 只传输已更改的数据
  • ✓ 使用 --delete 标志以避免重新下载已存在的文件
  • ✓ 先运行交易品种发现,确认准确、有效的货币对名称 —— 可避免因拼写错误导致同步失败
  • ✓ 保留日志文件以供审计
  • ✓ 为失败的下载设置 CloudWatch 告警
  • ✓ 在存储桶大小/文件数量发生重大变化时重新计算成本
  • ✓ 尽可能优先下载特定货币对,而非整个存档

解码与数据完整性

  • ✓ 解码前始终确认每个交易品种正确的点值(100,000 或 1,000)—— 错误的缩放会在无声无息中破坏价格数据
  • ✓ 将缺失的每日 .bi5 文件视为"当天没有逐笔数据"(周末/节假日),而非错误
  • ✓ 在未先检测格式的情况下,切勿在同一流水线中混用旧版按小时文件与新的按日文件
  • ✓ 验证 LZMA 解压模式 —— 某些库需要显式的"raw"模式(不含容器头)
  • ✓ 在导出前按时间顺序对解码后的逐笔数据进行排序,以防文件系统排序不一致

流水线与自动化

  • ✓ 使用 --cleanup 在解码后删除原始 .bi5 文件 —— 如有需要,随时可以从 S3 重新下载
  • ✓ 对于逐笔级别的数据,优先使用 Parquet 而非 CSV —— 通常体积小 5-10 倍且查询更快
  • ✓ 测试时使用 --start-date / --end-date 过滤器,以避免不必要的 S3 成本
  • ✓ 与交易品种发现相结合,自动遍历所有有效货币对,而非硬编码列表
  • ✓ 记录每次流水线运行(已下载/已解码/已导出的数量)以便审计

通用

  • ✓ 在非高峰时段运行下载以减少争用(而非降低成本)
  • ✓ 在关键场景中实施校验和以验证数据完整性
  • ✓ 随时间监控传输的数据总量,以控制累积成本