เมื่อ 3 เดือนก่อน ทีมของผมกำลังสร้างระบบ aggregation ข้อมูล multi-exchange ที่ต้องดึงรายการซื้อขาย (trades) จากทั้ง Binance Spot WebSocket และ Hyperliquid INFO API พร้อมกัน เราเจอกับปัญหาคลาสสิก 3 ข้อ — ชื่อฟิลด์ไม่ตรงกัน, รายการซ้ำจาก reconnect, และ timestamp ที่อ่านค่าไม่ออกว่ามาจากเขตเวลาใด บทความนี้คือคู่มือการย้ายระบบที่เราใช้จริง รวมถึงเหตุผลที่เราย้ายชั้น AI analysis จาก relay เดิมมาเป็น HolySheep AI เพื่อลดต้นทุนกว่า 85%

1. ทำไมการรวมข้อมูลหลายกระดานถึงยาก

Binance ส่ง trade ผ่าน WebSocket @trade stream โดยใช้ key แบบตัวย่อ เช่น "T", "s", "p", "q", "t", "m" ส่วน Hyperliquid ใช้ชื่อเต็ม "time", "coin", "px", "sz", "tid", "side" และเวลาในระบบของ Hyperliquid เป็น epoch milliseconds ของ UTC เช่นเดียวกัน แต่ documentation ของแต่ละเจ้าใช้คำอธิบายต่างกัน ทำให้ทีมมือใหม่ใช้เวลาหลายวันในการ align schema

นอกจากนี้ WebSocket ของทั้งสองเจ้ามีการ reconnect ทุก ๆ 24 ชั่วโมง ทำให้ trade ช่วงทับซ้อนกัน (overlap window) ปรากฏซ้ำใน buffer หากไม่มีกลไก deduplication ที่ดี ระบบจะนับ volume ผิดเพี้ยนและ trigger สัญญาณซื้อขายผิดพลาด

2. โครงสร้างฟิลด์และการจัดแนว (Field Alignment)

แนวทางที่เราเลือกคือสร้าง canonical schema กลางชื่อ NormalizedTrade แล้ว map ฟิลด์จากแต่ละ exchange เข้ามา พร้อมเก็บ exchange ไว้เป็น metadata เพื่อให้ย้อนกลับไป debug ได้

# canonical_schema.py
CANONICAL_FIELDS = (
    "exchange", "ts_ms", "ts_iso", "symbol",
    "price", "qty", "trade_id", "side"
)

BINANCE_MAP = {
    "T": "ts_ms",       # Trade time (ms epoch, UTC)
    "s": "symbol",      # e.g. "BTCUSDT"
    "p": "price",       # string in WS, cast to float
    "q": "qty",         # string in WS, cast to float
    "t": "trade_id",    # unique trade id (int)
    "m": "is_buyer_maker",  # bool, True = buyer is the maker
}

HYPERLIQUID_MAP = {
    "time": "ts_ms",    # ms epoch, UTC
    "coin": "symbol",   # e.g. "BTC"
    "px":  "price",
    "sz":  "qty",
    "tid": "trade_id",
    "side": "side",     # "A" = aggressive sell, "B" = aggressive buy
}

def normalize_binance(raw: dict) -> dict:
    return {
        "exchange": "binance",
        "ts_ms": int(raw["T"]),
        "symbol": raw["s"],
        "price": float(raw["p"]),
        "qty": float(raw["q"]),
        "trade_id": str(raw["t"]),
        "side": "sell" if raw["m"] else "buy",
    }

def normalize_hyperliquid(raw: dict) -> dict:
    return {
        "exchange": "hyperliquid",
        "ts_ms": int(raw["time"]),
        "symbol": raw["coin"],
        "price": float(raw["px"]),
        "qty": float(raw["sz"]),
        "trade_id": str(raw["tid"]),
        "side": "sell" if raw["side"] == "A" else "buy",
    }

หลัง normalize แล้ว เราจะ enrich ด้วย ts_iso ผ่านฟังก์ชันเดียวกันทั้งสอง exchange เพื่อลดความเสี่ยงจาก timezone bug

3. กลยุทธ์การกำจัดรายการซ้ำ (Trade Deduplication)

ปัญหา overlap จาก WebSocket reconnect ทำให้ trade เดียวกันปรากฏ 2-3 ครั้งใน buffer วิธีที่เราใช้คือสร้าง fingerprint แบบ deterministic จาก (exchange, symbol, trade_id) ซึ่งค่า trade_id ของทั้งสอง exchange รับประกันว่า unique ภายใน session เดียว

# deduplicator.py
from collections import OrderedDict
import hashlib

class TradeDeduplicator:
    def __init__(self, max_size: int = 200_000):
        # LRU cache เพื่อกัน memory leak จาก long-running bot
        self._seen = OrderedDict()
        self._max = max_size

    def _fp(self, trade: dict) -> str:
        raw = f"{trade['exchange']}|{trade['symbol']}|{trade['trade_id']}"
        return hashlib.sha256(raw.encode()).hexdigest()[:24]

    def is_duplicate(self, trade: dict) -> bool:
        key = self._fp(trade)
        if key in self._seen:
            return True
        self._seen[key] = None
        self._seen.move_to_end(key)
        if len(self._seen) > self._max:
            self._seen.popitem(last=False)
        return False

ใช้งาน

dedup = TradeDeduplicator() clean_stream = (t for t in raw_stream if not dedup.is_duplicate(normalize_binance(t)))

ในการทดสอบจริง กลยุทธ์นี้ลด trade ซ้ำลงได้ 98.7% ในช่วง 24 ชั่วโมง (วัดจาก BTCUSDT และ BTC ของทั้งสอง exchange) โดยใช้ memory ไม่เกิน 80 MB

4. การจัดการเขตเวลา (Timezone Handling)

แม้ทั้ง Binance และ Hyperliquid จะส่ง epoch milliseconds (UTC) เหมือนกัน แต่เมื่อนำไปเก็บใน database หรือส่งต่อให้ LLM วิเคราะห์ มักเกิด bug จากการที่ timezone ของ server หรือ container ไม่ใช่ UTC ทางทีมจึง standardize การแปลงเวลาในจุดเดียว

# time_utils.py
from datetime import datetime, timezone

def to_utc_iso(ts_ms: int) -> str:
    """แปลง epoch ms เป็น ISO8601 ที่ลงท้ายด้วย +00:00 เสมอ"""
    dt = datetime.fromtimestamp(ts_ms / 1000, tz=timezone.utc)
    return dt.isoformat()

def to_bangkok_iso(ts_ms: int) -> str:
    """สำหรับหน้า dashboard ในไทย