02 Model the data
ClickHouse Cloud에 타입이 지정된 market, tick, trade, 1분 집계 테이블을 만듭니다.
시작 지점
.env.polymarket이 source되어 있고, condition ID와 토큰 ID가 왜 다른지 알고 있습니다.
왜 이 테이블들인가
뒤에 나오는 다섯 개 쿼리는 관찰 대상 모든 시장에 대해 최근 시간 구간을 읽습니다. 그래서
이벤트 테이블의 키는 시간/시각 키로 시작하고, 그 다음에 그룹화에 쓰이는 토큰 또는 condition이
옵니다. 알려진 필드는 네이티브 타입을 사용합니다. 토큰 ID에는 UInt256, 이벤트 시각에는
DateTime64, 가격과 크기에는 정확한 십진수, 값의 범위가 제한된 이벤트 값에는 enum을
사용합니다. 불투명한 원본 페이로드는 어떤 쿼리도 그 필드를 읽지 않으므로 문자열로 둡니다.
파티션 키는 없습니다. 이 짧게 운영되는 워크샵에는 검증된 보존 경계가 없으므로, 수명 주기 요구사항이 생기기 전에 파티션을 추가하면 이득 없이 작은 파트만 만들게 됩니다.
Step 1 — 데이터베이스와 원본 테이블 생성
블록 전체를 ClickHouse Cloud SQL 콘솔에 복사해 실행하세요:
CREATE DATABASE IF NOT EXISTS polymarket;
CREATE TABLE IF NOT EXISTS polymarket.markets
(
market_id UInt64,
condition_id FixedString(66),
token_id UInt256,
outcome LowCardinality(String),
question String,
slug String,
active Bool,
accepting_orders Bool,
volume_24h Decimal128(8),
observed_at DateTime64(3, 'UTC')
)
ENGINE = ReplacingMergeTree(observed_at)
ORDER BY (condition_id, token_id);
CREATE TABLE IF NOT EXISTS polymarket.price_ticks
(
event_id FixedString(64),
condition_id FixedString(66),
token_id UInt256,
event_at DateTime64(3, 'UTC'),
observed_at DateTime64(3, 'UTC'),
event_kind Enum8(
'book_snapshot' = 1,
'price_change' = 2,
'last_trade_price' = 3,
'best_bid_ask' = 4,
'rest_book' = 5
),
source Enum8('WEBSOCKET' = 1, 'CLOB_REST' = 2, 'FIXTURE' = 3),
price Decimal64(12),
size Decimal128(8),
side Enum8('UNKNOWN' = 0, 'BUY' = 1, 'SELL' = 2),
best_bid Decimal64(12),
best_ask Decimal64(12),
midpoint Decimal64(12),
source_hash String,
raw_payload String
)
ENGINE = MergeTree
ORDER BY (toStartOfHour(event_at), token_id, event_at, event_id);
CREATE TABLE IF NOT EXISTS polymarket.trades
(
trade_id FixedString(64),
condition_id FixedString(66),
token_id UInt256,
event_at DateTime64(3, 'UTC'),
observed_at DateTime64(3, 'UTC'),
proxy_wallet FixedString(42),
side Enum8('UNKNOWN' = 0, 'BUY' = 1, 'SELL' = 2),
price Decimal64(12),
size Decimal128(8),
outcome LowCardinality(String),
transaction_hash FixedString(66),
title String
)
ENGINE = ReplacingMergeTree(observed_at)
ORDER BY (toStartOfHour(event_at), condition_id, event_at, trade_id);
CREATE OR REPLACE VIEW polymarket.trades_clean AS
SELECT *
FROM polymarket.trades FINAL;collector는 삽입 전에 중복을 방지합니다. ReplacingMergeTree는 두 번째 안전망입니다.
trades_clean 뷰는 병합이 아직 진행 중일 때도 작은 워크샵 쿼리들이 결정적으로 동작하게
해줍니다.
Step 2 — 1분 중간값 집계 생성
CREATE TABLE IF NOT EXISTS polymarket.market_midpoints_1m
(
token_id UInt256,
minute DateTime('UTC'),
open AggregateFunction(argMin, Decimal64(12), Tuple(DateTime64(3, 'UTC'), FixedString(64))),
high AggregateFunction(max, Decimal64(12)),
low AggregateFunction(min, Decimal64(12)),
close AggregateFunction(argMax, Decimal64(12), Tuple(DateTime64(3, 'UTC'), FixedString(64))),
updates AggregateFunction(count)
)
ENGINE = AggregatingMergeTree
ORDER BY (minute, token_id);
CREATE MATERIALIZED VIEW IF NOT EXISTS polymarket.market_midpoints_1m_mv
TO polymarket.market_midpoints_1m
AS
SELECT
token_id,
toStartOfMinute(event_at) AS minute,
argMinState(midpoint, tuple(event_at, event_id)) AS open,
maxState(midpoint) AS high,
minState(midpoint) AS low,
argMaxState(midpoint, tuple(event_at, event_id)) AS close,
countState() AS updates
FROM polymarket.price_ticks
WHERE midpoint > 0
AND event_kind IN ('book_snapshot', 'price_change', 'best_bid_ask', 'rest_book')
GROUP BY token_id, minute;이 materialized view는 호가 중간값만 집계합니다. 변경된 주문 단위 가격과 최종 거래 가격은 의도적으로 제외하므로 OHLC 시계열의 의미가 하나로 유지됩니다.
Step 3 — 모든 객체 검증
clickhouse client \
--host "$CLICKHOUSE_HOST" \
--port "$CLICKHOUSE_PORT" \
--user "$CLICKHOUSE_USER" \
--password "$CLICKHOUSE_PASSWORD" \
--secure \
--query "SHOW TABLES FROM polymarket"예상되는 이름에는 다음이 포함됩니다:
market_midpoints_1m
market_midpoints_1m_mv
markets
price_ticks
trades
trades_clean완료 조건
로컬 ClickHouse 서버를 실행하지 않은 상태에서 SHOW TABLES가 여섯 개 객체를 모두 반환합니다.
다음: 실시간 collector를 시작합니다.