03 流式接入实时数据
运行具备韧性的采集器,验证 WebSocket、REST 对账和写入 Cloud 的过程。
macOS terminal: Run workshop commands in Terminal using zsh or bash.
起点
六个 polymarket 对象都已存在,并且已 source .env.polymarket。
启动了什么
一个无状态的 Python 容器:
- 通过 Gamma 发现五个活跃市场;
- 在公开的 CLOB WebSocket 上订阅两个结果选项的 token;
- 每 10 秒对账一次公开交易;
- 在 WebSocket 停滞时轮询 CLOB 盘口;以及
- 用带确认的异步插入把数据写入 ClickHouse Cloud。
这里没有本地数据库、消息代理、仪表板服务器,也不需要任何 Polymarket 凭据。
第 1 步:构建并启动采集器
docker compose --env-file .env.polymarket up -d --build collector
docker compose --env-file .env.polymarket ps在完成市场发现和第一次成功写入 Cloud 之后,状态会变为 healthy。当 REST 数据保持
最新、WebSocket 正在重连时,应用层状态为 degraded 仍然算 Docker 健康。
第 2 步:读懂健康检查约定
curl --fail --silent http://localhost:8090/health \
| python3 -m json.tool预期字段:
{
"status": "live",
"websocket": "connected",
"queue_depth": 0,
"queue_capacity": 10000,
"watched_markets": 5,
"watched_tokens": 10,
"fresh_tokens": 10
}如果 last_trade_reconcile_at 和 last_book_fallback_at 在持续推进,那么
status: degraded 并带有 reason: websocket_stale_rest_active 是可以接受的。unhealthy 则
不可接受;请查看故障排查页面。
第 3 步:观察数据源和写入事件
docker compose --env-file .env.polymarket logs --tail=30 collector日志是 JSON 格式。找到 collector_ready。数据源或 ClickHouse 失败时会附带
一段有长度限制的错误预览和重试延迟;日志中不会记录密码。
第 4 步:证明数据行已进入 Cloud
clickhouse client \
--host "$CLICKHOUSE_HOST" \
--port "$CLICKHOUSE_PORT" \
--user "$CLICKHOUSE_USER" \
--password "$CLICKHOUSE_PASSWORD" \
--secure \
--query "
SELECT 'markets' AS table, count() AS rows FROM polymarket.markets
UNION ALL
SELECT 'quote_midpoints', countIf(midpoint > 0) FROM polymarket.price_ticks
UNION ALL
SELECT 'trades', count() FROM polymarket.trades_clean
UNION ALL
SELECT 'one_minute_states', count() FROM polymarket.market_midpoints_1m
"在进入模块 04 之前,markets、quote_midpoints 和 one_minute_states 必须都大于
零。trades 通常在一分钟内就开始增长;交易清淡的市场可能会更慢。
第 5 步:只在必要时使用确定性的 fixture 模式
如果场地网络屏蔽了 Polymarket,或者所选市场在 60 秒后仍没有任何变动:
sed -i.bak 's/^POLYMARKET_MODE=.*/POLYMARKET_MODE=fixture/' .env.polymarket
set -a; source ./.env.polymarket; set +a
docker compose --env-file .env.polymarket up -d --build --force-recreate collector再次运行健康检查和行数检查。预期状态:fixture;tick 和交易
计数每五秒增长一次。在本模块结束前请保留备份文件。
完成标准
- 健康状态为
live、带有最新 REST 时间戳的degraded,或者fixture; watched_markets为 5;以及- 四个 Cloud 行数都返回了结果,且报价行和一分钟行都大于零。
下一步:查询增量聚合。