04 Aggregate เรียลไทม์
อ่าน OHLC ของ quote-midpoint รายหนึ่งนาทีที่ดูแลโดย materialized view แบบเพิ่มทีละส่วน
จุดเริ่มต้น
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เป็นค่าปัจจุบันเมื่อฟีดยังทำงานอยู่
ต่อไป: สืบสวนการเคลื่อนไหวของตลาด