04 Aggregate thời gian thực
Đọc OHLC quote-midpoint theo từng phút được duy trì bởi một materialized view tăng dần.
Điểm khởi đầu
Collector đang healthy và polymarket.price_ticks đã có dữ liệu.
Tại sao
Người dùng dashboard liên tục hỏi cùng một chuỗi theo phút. Tính nó một lần khi các block mới tới sẽ chuyển công việc từ mỗi lần làm mới dashboard sang thời điểm insert. Bảng thô vẫn sẵn sàng cho các câu hỏi tùy biến.
Bước 1 — Truy vấn các aggregate state cho đúng cách
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 được view ghi vào; truy vấn hoàn tất chúng bằng các hàm Merge
tương ứng. Open và close dùng thời điểm sự kiện cộng với event ID tất định,
nhờ đó các bản ghi đến không đúng thứ tự và các trường hợp trùng cùng một millisecond đều được
giải quyết một cách nhất quán.
Bước 2 — So sánh số dòng đọc giữa truy vấn thô và truy vấn aggregate
Chạy phiên bản thô tương đương:
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;Trong SQL console, hãy so sánh read rows của cả hai truy vấn. Truy vấn aggregate đọc các state ở
mức block mà AggregatingMergeTree gộp lại ở chế độ nền, thường là ít dòng hơn rất nhiều
so với việc quét từng bản cập nhật nguồn.
Bước 3 — Kiểm chứng rằng view được dẫn dắt bởi 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;Ở chế độ live hoặc fixture, newest_minute tiến lên mà không cần job làm mới theo lịch nào.
Hoàn thành khi
- truy vấn aggregate trả về các dòng OHLC;
- open/high/low/close là các xác suất từ 0 đến 100; và
newest_minutelà hiện thời với một feed đang hoạt động.
Tiếp theo: điều tra một biến động thị trường.