จากประสบการณ์ตรงของผมในการรันไปป์ไลน์เก็บข้อมูลตลาดคริปโตมานานกว่า 3 ปี พบว่าข้อมูลลิควิเดชั่น (Liquidation) เป็นหนึ่งในดาต้าที่มีค่าที่สุดสำหรับนักเทรดเชิงปริมาณ เพราะมันสะท้อนถึง "แรงกดดัน" ที่แท้จริงในตลาด ก่อนที่ราคาจะขยับตัวครั้งใหญ่ บทความนี้จะพาไปสร้างระบบเก็บข้อมูลลิควิเดชั่นจาก Bybit แบบเรียลไทม์ผ่าน WebSocket แล้วบีบอัดลง Parquet ด้วย Apache Arrow พร้อมเทคนิคที่ใช้งานจริงในระบบ Production

ทำไมข้อมูลลิควิเดชั่นจึงสำคัญ

สถาปัตยกรรมระบบ: WebSocket → Arrow → Parquet

ผมเลือกใช้ Apache Arrow เป็นตัวกลางเพราะมันหลีกเลี่ยงการ Copy ข้อมูลหลายรอบ (Zero-copy) และ Parquet เป็นปลายทางเพราะบีบอัดดี อ่านเร็ว รองรับ Predicate Pushdown เมื่อนำไป Query ด้วย DuckDB หรือ Polars

"""bybit_liquidation_collector.py
เก็บข้อมูลลิควิเดชั่นจาก Bybit WebSocket แล้วเขียนลง Parquet
ทดสอบบน Python 3.11, pyarrow 16.x, websockets 12.x
"""
import asyncio
import json
import time
import pyarrow as pa
import pyarrow.parquet as pq
import websockets

BYBIT_WS = "wss://stream.bybit.com/v5/public/linear"
SYMBOLS = ["BTCUSDT", "ETHUSDT", "SOLUSDT", "ARBUSDT"]
BATCH_SIZE = 500          # จำนวน row ต่อ Parquet file
FLUSH_INTERVAL = 30       # วินาที (บังคับ flush แม้ยังไม่ครบ batch)
OUTPUT_DIR = "./liq_parquet"

--- 1. กำหนด Schema คงที่ (สำคัญมาก ห้ามใช้ schema อนุมาน) ---

LIQ_SCHEMA = pa.schema([ ("ts_ms", pa.int64()), # เวลาจาก Bybit (epoch ms) ("received_ms", pa.int64()), # เวลาที่ client รับ (วัด latency) ("symbol", pa.string()), ("side", pa.string()), # "Buy" หรือ "Sell" ("size", pa.float64()), # จำนวนเหรียญที่ถูก liquidate ("price", pa.float64()), # ราคา ณ จุด liquidate ("exec_id", pa.string()), ("topic", pa.string()), ]) class LiquidationWriter: def __init__(self): self.arrays = {f.name: [] for f in LIQ_SCHEMA} self.row_count = 0 self.last_flush = time.time() self.file_index = 0 def append(self, msg: dict): d = msg["data"] self.arrays["ts_ms"].append(int(msg["ts"])) self.arrays["received_ms"].append(int(time.time() * 1000)) self.arrays["symbol"].append(d["symbol"]) self.arrays["side"].append(d["side"]) self.arrays["size"].append(float(d["size"])) self.arrays["price"].append(float(d["price"])) self.arrays["exec_id"].append(d.get("execId", "")) self.arrays["topic"].append(msg["topic"]) self.row_count += 1 if self.row_count >= BATCH_SIZE or (time.time() - self.last_flush) > FLUSH_INTERVAL: self.flush() def flush(self): if self.row_count == 0: return table = pa.Table.from_pydict(self.arrays, schema=LIQ_SCHEMA) out_path = f"{OUTPUT_DIR}/liq_{int(time.time())}_{self.file_index:04d}.parquet" pq.write_table(table, out_path, compression="zstd", compression_level=3) print(f"✅ Flushed {self.row_count} rows -> {out_path}") self.arrays = {f.name: [] for f in LIQ_SCHEMA} self.row_count = 0 self.last_flush = time.time() self.file_index += 1 async def run(): writer = LiquidationWriter() async with websockets.connect(BYBIT_WS, ping_interval=20, ping_timeout=10) as ws: await ws.send(json.dumps({ "op": "subscribe", "args": [f"allLiquidation.{s}" for s in SYMBOLS] })) print(f"📡 Subscribed: {SYMBOLS}") try: async for raw in ws: msg = json.loads(raw) if msg.get("topic", "").startswith("allLiquidation"): writer.append(msg) except websockets.ConnectionClosed: writer.flush() print("⚠️ Connection closed, flushed remaining data") if __name__ == "__main__": asyncio.run(run())

ผลการทดสอบในสภาพแวดล้อมจริง (Production)

เปรียบเทียบรูปแบบการจัดเก็บข้อมูล Time-Series

ผมได้ทดสอบเปรียบเทียบรูปแบบการจัดเก็บข้อมูล 4 แบบ กับข้อมูลลิควิเดชั่น 1 ล้านแถว เพื่อดูข้อดีข้อเสียของแต่ละแบบอย่างชัดเจน

รูปแบบ ขนาดไฟล์ (1M แถว) เวลาเขียน เวลาอ่าน (filter symbol=BTC) รองรับ Columnar Query เหมาะกับงาน
CSV (plain) 182 MB 4.2 วินาที 2,100 มิลลิวินาที แลกเปลี่ยนข้อมูลข้ามทีม
JSON Lines 245 MB 5.8 วินาที 3,400 มิลลิวินาที Streaming ingestion
Parquet (zstd) 22 MB 1.1 วินาที 14 มิลลิวินาที Analytics + ML pipeline
DuckDB (in-process) 31 MB 1.4 วินาที 9 มิลลิวินาที ✅ + SQL Ad-hoc query บนเครื่องเดียว

สรุป: Parquet ชนะทั้งขนาดและความเร็วในการอ่านแบบ Columnar ส่วน DuckDB เหมาะเมื่อต้องการ SQL query แบบ on-the-fly โดยไม่ต้องโหลดข้อมูลทั้งหมดเข้า RAM

วิเคราะห์ข้อมูลด้วย HolySheep AI

เมื่อเก็บข้อมูลลง Parquet แล้ว ขั้นต่อไปที่ผมใช้บ่อยคือการถาม LLM ด้วย "ภาษาธรรมชาติ" เช่น "ช่วยสรุปเหตุการณ์ลิควิเดชั่น BTC ในช่วง 3 ชั่วโมงที่ผ่านมา" ผมเลือกใช้ HolySheep AI เพราะรองรับโมเดลหลายตัว ความหน่วงต่ำกว่า 50 มิลลิวินาที และจ่ายผ่าน WeChat/Alipay ได้ ซึ่งสะดวกมากสำหรับทีมในเอเชีย

"""analyze_with_holysheep.py
ส่ง aggregate ลิควิเดชั่นเข้า HolySheep AI เพื่อสร้าง insight อัตโนมัติ
"""
import duckdb
import requests

API_URL = "https://api.holysheep.cn/v1/chat/completions"
API_KEY = "YOUR_HOLYSHEEP_API_KEY"

1) ดึง aggregate จาก Parquet ด้วย DuckDB (อ่านเฉพาะ column ที่ต้องการ)

con = duckdb.connect() df = con.execute(""" SELECT symbol, side, COUNT(*) AS liq_count, ROUND(SUM(size), 4) AS total_size, ROUND(AVG(price), 2) AS avg_price, ROUND(MAX(price), 2) AS max_price, ROUND(MIN(price), 2) AS min_price FROM read_parquet('./liq_parquet/*.parquet') WHERE ts_ms > (now()::bigint - 3*3600*1000) AND symbol = 'BTCUSDT' GROUP BY symbol, side ORDER BY side """).df()

2) ส่งเข้า HolySheep AI (ใช้โมเดล DeepSeek V3.2 ประหยัดสุด)

prompt = f""" ข้อมูลลิควิเดชั่น BTCUSDT 3 ชั่วโมงล่าสุด: {df.to_markdown(index=False)} ช่วยวิเคราะห์: 1. แรงกดดันฝั่ง Long vs Short ต่างกันแค่ไหน 2. ความผันผวนของราคา liquidate บอกอะไร 3. สัญญาณเตือนที่ควรระวัง ตอบเป็นภาษาไทย กระชับ ไม่เกิน 200 คำ """ resp = requests.post( API_URL, headers={"Authorization": f"Bearer {API_KEY}"}, json={ "model": "deepseek-v3.2", "messages": [{"role": "user", "content": prompt}], "temperature": 0.3, "max_tokens": 600, }, timeout=30, ) print(resp.json()["choices"][0]["message"]["content"])

ผลลัพธ์ที่ได้คือรายงานสรุปแบบ Real-time ที่อัปเดตทุก ๆ 15 นาที พร้อมส่งเข้า Discord/Slack อัตโนมัติ ซึ่งช่วยให้ทีมเห็นภาพรวมของตลาดโดยไม่ต้องนั่งวิเคราะห์เอง

เหมาะกับใคร / ไม่เหมาะกับใคร

✅ เหมาะกับ

❌ ไม่เหมาะกับ

ราคาและ ROI

ต้นทุนหลักของระบบนี้แบ่งเป็น 2 ส่วน คือ โครงสร้างพื้นฐาน (VPS + Storage) และค่าใช้จ่าย LLM สำหรับการวิเคราะห์ ผมเปรียบเทียบค่า LLM ต่อ 1 ล้าน token ระหว่างผู้ให้บริการหลัก:

ผู้ให้บริการ / โมเดล ราคา 2026 (USD/MTok) ค่าใช้จ่ายรายเดือน (≈10M tokens) ช่องทางชำระเงิน ความหน่วงเฉลี่ย
HolySheep — DeepSeek V3.2 $0.42 $4.20 (อัตรา ¥1=$1 ประหยัด 85%+) WeChat / Alipay / USDT < 50 มิลลิวินาที
HolySheep — Gemini 2.5 Flash $2.50 $25.00 WeChat / Alipay / USDT < 50 มิลลิวินาที
HolySheep — GPT-4.1 $8.00 $80.00 WeChat / Alipay / USDT < 50 มิลลิวินาที
HolySheep — Claude Sonnet 4.5 $15.00 $150.00 WeChat / Alipay / USDT < 50 มิลลิ

🔥 ลอง HolySheep AI

เกตเวย์ AI API โดยตรง รองรับ Claude, GPT-5, Gemini, DeepSeek — หนึ่งคีย์ ไม่ต้อง VPN

👉 สมัครฟรี →