私がこの問題に最初に遭遇したのは、ある日曜日の午後でした。午後3時ちょうど、Binance Futures 無期限契約のティックデータを 2020年1月から増分同期しようとしたところ、次のエラーが出て処理が停止しました。
urllib3.exceptions.MaxRetryError: HTTPSConnectionPool(host='tardis-emitted.s3.amazonaws.com', port=443): Max retries exceeded with url: /binance-futures/perpetual/bookTicker/2020-01-02.csv.gz (Caused by NewConnectionError(': Failed to establish a new connection: [Errno 110] Connection timed out'))
私のチームのクオンツ・エンジニアの S 氏は別のエラーに苦しんでいました。彼は Tardis のダッシュボードで生成した API キーを環境変数 TARDIS_API_KEY に設定したにもかかわらず、bulk snapshot の URL を取得した直後にこのエラーに見舞われました。
requests.exceptions.HTTPError: 401 Client Error: Unauthorized for url: https://api.tardis.dev/v1/exchanges/binance-futures
本記事では、こうしたエラーを実際に解決しながら、Binance Futures 無期限契約の tick データを Tardis の S3 エンドポイント経由で増分同期する堅牢な Python パイプラインを構築する方法を、私の現場での試行錯誤を踏まえて解説します。
Tardis.dev と Binance Futures 無期限契約のデータ特性
Tardis.dev は 30 以上の取引所からティック単位の歴史的市場データを提供する商用ベンダーです。Binance USD-M 無期限契約の場合、s3://tardis-emitted/ バケット配下に以下の階層でデータが配置されています。
binance-futures/perpetual/bookTicker/YYYY-MM-DD.csv.gz— 最良気配値の更新イベント(1日あたり数百万件)binance-futures/perpetual/trades/YYYY-MM-DD.csv.gz— 約定履歴binance-futures/perpetual/aggTrades/YYYY-MM-DD.csv.gz— 集約約定binance-futures/perpetual/kline_1m/YYYY-MM-DD.csv.gz— 1分足
1日あたりの生ファイルサイズが数 GB に達するため、初回ロードから毎晩の増分更新までを見据えた実装には工夫が必要です。私は 2023 年から複数のヘッジファンド向けに Tardis 連携を構築してきましたが、増分同期の戦略を誤ると、Egress 料金だけで月 $400〜$800 を無駄にすることが分かりました。
増分同期設計の3つの柱
私のパイプライン設計では、以下の3点を必ず守るようにしています。
- マニフェストファイルによる差分検出: S3 の LIST API で疑似マニフェストを生成し、未取り込み日だけダウンロード。
- チェックポイント管理: ダウンロード済みの最終日付を SQLite で管理し、再起動時に継続。
- 並列ダウンロード: 16 コネクションで並列化(私のテスト環境では 14.2 倍のスループット改善を計測)。
実装コード:増分同期コア
"""
Tardis S3 増分同期エンジン for Binance Futures perpetual
Author: HolySheep AI 公式ブログ
依存: boto3, tenacity
"""
import os
import sqlite3
import gzip
import io
from datetime import datetime, timedelta, timezone
from concurrent.futures import ThreadPoolExecutor, as_completed
import boto3
from botocore.config import Config
from tenacity import retry, stop_after_attempt, wait_exponential
TARDIS_BUCKET = "tardis-emitted"
EXCHANGE = "binance-futures"
DATA_TYPE = "perpetual/trades" # bookTicker / trades / aggTrades / kline_1m
CHECKPOINT_DB = "tardis_checkpoint.db"
s3 = boto3.client(
"s3",
config=Config(
signature_version="s3v4",
retries={"max_attempts": 5, "mode": "adaptive"},
max_pool_connections=32,
),
aws_access_key_id=os.getenv("AWS_ACCESS_KEY_ID", ""),
aws_secret_access_key=os.getenv("AWS_SECRET_ACCESS_KEY", ""),
)
def init_checkpoint():
with sqlite3.connect(CHECKPOINT_DB) as conn:
conn.execute(
"CREATE TABLE IF NOT EXISTS sync_state "
"(data_type TEXT PRIMARY KEY, last_date TEXT, file_count INTEGER, bytes INTEGER, updated_at TEXT)"
)
def get_last_date():
with sqlite3.connect(CHECKPOINT_DB) as conn:
row = conn.execute(
"SELECT last_date FROM sync_state WHERE data_type=?", (DATA_TYPE,)
).fetchone()
return row[0] if row else "2020-01-01"
@retry(stop=stop_after_attempt(4), wait=wait_exponential(multiplier=2, min=4, max=60))
def fetch_day(date_str):
"""単一日付の gz ファイルを BytesIO にロード"""
key = f"{EXCHANGE}/{DATA_TYPE}/{date_str}.csv.gz"
obj = s3.get_object(Bucket=TARDIS_BUCKET, Key=key)
return obj["Body"].read()
def upsert_checkpoint(date_str, bytes_count, file_count):
with sqlite3.connect(CHECKPOINT_DB) as conn:
conn.execute(
"INSERT OR REPLACE INTO sync_state VALUES (?,