PolymarketClickHouse Workshops

04 Aggregate เรียลไทม์

อ่าน OHLC ของ quote-midpoint รายหนึ่งนาทีที่ดูแลโดย materialized view แบบเพิ่มทีละส่วน

Your computer
macOS terminal: Run workshop commands in Terminal using zsh or bash.

จุดเริ่มต้น

collector อยู่ในสถานะ healthy และ polymarket.price_ticks มีแถวข้อมูลแล้ว

ทำไป

ผู้ใช้แดชบอร์ดขอชุดข้อมูลรายหนึ่งนาทีชุดเดิมซ้ำ ๆ การคำนวณครั้งเดียวตอนที่บล็อกใหม่ เข้ามาย้ายงานจากการรีเฟรชแดชบอร์ดทุกครั้งไปไว้ที่ตอน insert ตารางข้อมูลดิบ ยังคงใช้ได้สำหรับคำถามเฉพาะกิจ

ขั้นที่ 1 — คิวรี aggregate state อย่างถูกวิธี

SELECT
    minute,
    token_id,
    round(argMinMerge(open) * 100, 2) AS open_percent,
    round(maxMerge(high) * 100, 2) AS high_percent,
    round(minMerge(low) * 100, 2) AS low_percent,
    round(argMaxMerge(close) * 100, 2) AS close_percent,
    countMerge(updates) AS updates
FROM polymarket.market_midpoints_1m
WHERE minute >= now() - INTERVAL 30 MINUTE
GROUP BY minute, token_id
ORDER BY minute DESC, token_id
LIMIT 30;

argMinState/argMaxState ถูกเขียนโดย view ตัวคิวรีทำให้ค่าเป็นผลสุดท้ายด้วยฟังก์ชัน Merge ที่เข้าคู่กัน open และ close ใช้เวลาเหตุการณ์บวกกับ event ID แบบ deterministic ดังนั้นข้อมูลที่มาถึงไม่เป็นลำดับและค่าที่เสมอกันในมิลลิวินาทีเดียวกันจะลงตัวอย่างสอดคล้องกัน

ขั้นที่ 2 — เทียบจำนวนแถวที่คิวรีข้อมูลดิบและคิวรี aggregate อ่าน

รันคิวรีที่เทียบเท่ากันบนข้อมูลดิบ:

SELECT
    toStartOfMinute(event_at) AS minute,
    token_id,
    round(argMin(midpoint, tuple(event_at, event_id)) * 100, 2) AS open_percent,
    round(max(midpoint) * 100, 2) AS high_percent,
    round(min(midpoint) * 100, 2) AS low_percent,
    round(argMax(midpoint, tuple(event_at, event_id)) * 100, 2) AS close_percent,
    count() AS updates
FROM polymarket.price_ticks
WHERE midpoint > 0
  AND event_at >= now() - INTERVAL 30 MINUTE
  AND event_kind IN ('book_snapshot', 'price_change', 'best_bid_ask', 'rest_book')
GROUP BY minute, token_id
ORDER BY minute DESC, token_id
LIMIT 30;

ในคอนโซล SQL ให้เทียบ read rows ของทั้งสองคิวรี ตัว aggregate อ่าน state ระดับบล็อก ที่ AggregatingMergeTree รวมให้เบื้องหลัง ซึ่งปกติมีจำนวนแถวน้อยกว่า การสแกนทุกการอัปเดตต้นทางอย่างมาก

ขั้นที่ 3 — ตรวจสอบว่า view ทำงานตามการ insert

SELECT
    max(minute) AS newest_minute,
    dateDiff('second', newest_minute, now()) AS age_seconds,
    countMerge(updates) AS source_updates
FROM polymarket.market_midpoints_1m;

ในโหมดสดหรือโหมด fixture ค่า newest_minute เดินหน้าต่อโดยไม่ต้องมีงานรีเฟรชตามตาราง

ถือว่าเสร็จเมื่อ

  • คิวรี aggregate คืนแถว OHLC ออกมา
  • open/high/low/close เป็นความน่าจะเป็นระหว่าง 0 ถึง 100 และ
  • newest_minute เป็นค่าปัจจุบันเมื่อฟีดยังทำงานอยู่

ต่อไป: สืบสวนการเคลื่อนไหวของตลาด

ในหน้านี้

Track your progress?

Optional. We email a link to confirm your address; progress records once you open it.

Please use your work email address, not a personal one.

Progress tracking also requires accepting the current Terms of Service in Privacy settings.

TH