04 实时聚合
读取由增量 materialized view 维护的一分钟报价中间价 OHLC。
macOS terminal: Run workshop commands in Terminal using zsh or bash.
起点
采集器处于健康状态,并且 polymarket.price_ticks 中已有数据行。
为什么
仪表板用户会反复请求同一个一分钟序列。在新数据块到达时计算一次, 就把工作量从每次仪表板刷新转移到了插入时刻。原始表 仍然可用于临时的探索性问题。
第 1 步:正确查询聚合状态
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 是由视图写入的;查询用对应的
Merge 函数将它们最终化。open 和 close 使用事件时间加上确定性的 event ID,
因此乱序到达和同一毫秒内的并列都能得到一致的结果。
第 2 步:对比原始查询和聚合查询读取的行数
运行原始表上的等价查询:
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。聚合查询读取的是
AggregatingMergeTree 在后台合并的块级状态,通常远少于
扫描每一条源更新的行数。
第 3 步:验证视图是由插入驱动的
SELECT
max(minute) AS newest_minute,
dateDiff('second', newest_minute, now()) AS age_seconds,
countMerge(updates) AS source_updates
FROM polymarket.market_midpoints_1m;在 live 或 fixture 模式下,newest_minute 无需任何定时刷新任务就会持续推进。
完成标准
- 聚合查询返回了 OHLC 数据行;
- open/high/low/close 是介于 0 到 100 之间的概率值;以及
- 在数据源活跃的情况下,
newest_minute是最新的。
下一步:调查一次市场变动。