
周二的体育App跑得稳稳当当,什么问题都没有。
然后周六下午三点到了。几十场比赛在同一分钟内开球。进球同时涌进好几个联赛的数据流。十分钟后,所有用户都在疯狂刷新,而你的单线程轮询器正在疯狂轰炸那个刚刚返回429的API。
所有"测试时没问题"的东西,其实都建立在一个前提上:流量很乖。
这篇教程要搞定的是:一个在流量不乖的时候也能正常运行的体育数据消费者。用Python和Redis来搭,而且每个组件都可以复用到任何实时数据源上,不只是体育。
学完之后你会得到:
一个毫秒级响应、签名验证的Webhook接收器
一个能吸收突发流量而不丢数据的持久化队列(Redis Streams)
能扛住重复、重试和乱序事件的幂等Worker
支持共享预算、退避、抖动和Retry-After处理的限速API客户端
单飞(single-flight)和过期重刷(stale-while-revalidate)的旁路缓存,一千个用户只产生一次API调用
能修复Webhook遗漏内容的对账轮询器
能验证以上所有功能能否正常工作的流量模拟器
开始搞。
为什么比赛日会让简陋的消费者挂掉
体育流量不是平滑的,它是同步的,而同步正是杀死系统的东西。
开球时间是同步的。英超周六比赛日,一整轮比赛都在同一时刻开球。安静了一周的比赛突然全部变成直播。
事件是同步的。一堆进球、红牌、换人在同一秒内到达跨多场比赛,然后安静一分钟。
终场哨声是同步的。一起开球的比赛一起结束,所以一大波match.finished事件几乎在同一时刻涌入。
用户也是同步的。大比赛进球时,数万人同时去刷新页面。
运动项目重叠。板球T20决赛、网球比赛日和足球比赛日可能同时达到峰值,所以峰值不限于一个运动项目。
简陋的消费者会以可预测的方式挂掉:
简陋做法 比赛日会出现什么问题
每个用户请求发一次API调用 你的用户数量直接变成你的限速问题
每隔几秒轮询每场比赛 请求数量随比赛数量增长,而不是随价值增长
在HTTP处理器里同步处理事件 一次慢的数据库写入就会让发送方超时
信任到达顺序 重试的"进球"事件在"终场比分"之后到达,把比赛状态改回直播
收到429立即重试 你把限速问题变成了一场雷暴
不做去重 重试导致进球被重复计数
每一条都有简单、成熟的解决方案。我们来一个一个加上。
架构设计
整个设计的黄金法则:用户请求永远不应该触发对上游API的调用。用户从你的缓存读数据。上游API的喂数据由你的Worker控制,速度由你决定。
选择数据接收方式
在写代码之前,先决定数据怎么进入你的系统。实时比分API有三种方式推送相同的事件:
方式 适合场景 权衡
REST轮询 简单任务、对账 成本随轮询频率增长;总有延迟
WebSocket 延迟敏感的实时界面 你需要自己处理重连逻辑和连接限制
Webhook 事件驱动的后端 需要一个公开的HTTPS端点,快速响应
对于比赛日规模,我建议Webhook作为主路径,REST作为安全垫:
Webhook是推送模式,所以请求数量不会随直播比赛数量增长。
它自带健壮消费者需要的属性:签名负载、自动重试带退避、用于去重的event_id、用于排序的每比赛序列号。
如果你的端点宕机一段时间,交付记录可以被列出并重放。
也会展示WebSocket方案,因为很多团队更喜欢它。
写代码前先做容量计算
先用纸笔算一下。这些数字是举例假设,不是真实API的限制:
简陋方案:每10秒轮询每场直播比赛
300场直播比赛 ÷ 10秒 = 30请求/秒 = 1800请求/分钟
好一点:每10秒每个运动项目发一次"直播"批量调用
13个运动项目 ÷ 10秒 ≈ 78请求/分钟
最好:Webhook作为主路径,每30秒轮询一次作为安全垫
13个运动项目 ÷ 30秒 ≈ 26请求/分钟
你的用户:
50000用户每10秒刷新一次 = 5000读/秒
→ 这些请求必须打到你的缓存,不能打到上游API
同样数据,用烂方法要1800请求/分钟,用好方法只要26请求/分钟。不同套餐的计划限制不同,先看价格页面和真实响应的限速头,然后用那个数字来设置BUDGET_PER_MIN。
第一步:初始化
在注册页面申请API密钥。
阅读文档了解认证方式和约定,然后在沙箱里试一下端点,看看真实负载长什么样再动手。
把API参考文档开着。它记录了X-RateLimit-*头、429和Retry-After行为,以及标准响应格式(data、meta、errors)。
项目结构:
matchday/
├── .env
├── docker-compose.yml
├── config.py
├── app.py # webhook接收器 + 读API + SSE
├── worker.py # 流消费者、幂等状态更新
├── client.py # 限速、带缓存的REST客户端
├── poller.py # 对账安全垫
├── ws_consumer.py # 可选的WebSocket摄入
├── register.py # 注册webhook
└── simulate.py # 流量模拟器 + 验证器
安装依赖:
python -m venv .venv && source .venv/bin/activate
pip install fastapi uvicorn "redis>=5" requests httpx websockets python-dotenv docker-compose.yml(Redis Streams需要Redis 7):
services: redis: image: redis:7-alpine command: ["redis-server", "--appendonly", "yes"] ports: ["6379:6379"] volumes: ["redis-data:/data"] volumes: redis-data: .env:
ORBISTATS_API_KEY=your_key_here
WEBHOOK_SECRET=generate_a_long_random_string
REDIS_URL=redis://localhost:6379/0 # Stay below your plan's real limit (aim for ~70-80%)
BUDGET_PER_MIN=60 # Docs show two live routes; confirm in the sandbox and set one:
# {sport}/matches/live (API reference)
# {sport}/live (Live Scores page)
LIVE_PATH={sport}/matches/live 生成一个强随机webhook密钥:
python -c "import secrets; print(secrets.token_urlsafe(32))" 启动Redis:
docker compose up -d 第二步:共享配置
config.py:
import os
from dotenv import load_dotenv load_dotenv() API_KEY = os.environ["ORBISTATS_API_KEY"]
WEBHOOK_SECRET = os.environ["WEBHOOK_SECRET"]
REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/0")
BASE_URL = "https://api.orbistats.com/v1" BUDGET_PER_MIN = int(os.getenv("BUDGET_PER_MIN", "60"))
LIVE_PATH = os.getenv("LIVE_PATH", "{sport}/matches/live") STREAM = "events" # durable queue of incoming events
GROUP = "workers" # consumer group
DLQ = "events:dead" # events that failed repeatedly 第三步:Webhook接收器(快速响应,其他什么都不做)
Webhook文档里有一条规则:5秒内必须响应,否则这次交付算失败并重试。经验很简单:在HTTP处理器里做最少的事。验签、去重、入队、返回200。所有慢的操作都在后面、在Worker里做。
每次交付都是一个带签名的POST。签名是用你的密钥对原始请求体计算的HMAC-SHA256,通过X-Orbistats-Signature头发送。文档里的负载格式:
{ "event": "match.finished", "event_id": "evt_9f2a1c4b", "delivery_id": "dlv_7c3e08a1", "sequence": 214, "match_id": 48213, "sport": "football", "final_score": "2-1", "timestamp": "2026-09-14T16:52:11Z"
} 两个字段对扩展性最关键:
event_id在重试时保持稳定,用来去重。
sequence对每个match_id单调递增,用来排序,而不是按到达时间。
app.py:
import hashlib
import hmac
import json import redis.asyncio as aioredis
from fastapi import FastAPI, HTTPException, Request, Response
from fastapi.responses import StreamingResponse from config import DLQ, GROUP, REDIS_URL, STREAM, WEBHOOK_SECRET app = FastAPI()
r = aioredis.from_url(REDIS_URL, decode_responses=True) def valid_signature(raw: bytes, header: str) -> bool: if not header.startswith("sha256="): return False digest = hmac.new(WEBHOOK_SECRET.encode(), raw, hashlib.sha256).hexdigest() # constant-time comparison avoids timing attacks return hmac.compare_digest(f"sha256={digest}", header) @app.post("/webhook")
async def webhook(request: Request): raw = await request.body() # sign the RAW bytes, never re-serialized JSON if not valid_signature(raw, request.headers.get("X-Orbistats-Signature", "")): raise HTTPException(status_code=401, detail="bad signature") try: event = json.loads(raw) event_id = event["event_id"] except (ValueError, KeyError): raise HTTPException(status_code=400, detail="malformed event") # First writer wins. The key outlives the longest retry window (~1h) by a lot. key = f"seen:{event_id}" first_time = await r.set(key, 1, nx=True, ex=86400) if first_time: try: await r.xadd(STREAM, {"payload": raw.decode()}, maxlen=500_000, approximate=True) except Exception: # If we failed to enqueue, forget we saw it so the sender's retry works. await r.delete(key) raise HTTPException(status_code=503, detail="queue unavailable") return Response(status_code=200) # duplicates are acknowledged too 三个值得抄走的细节:
验签在解析之前做。未经认证的负载不值得信任。
> 重复请求返回200。发送方应该停止重试,虽然我们忽略了这个事件。
> "失败时忘记它"的分支。没有这个分支,入队失败会标记事件为已看到,重试就会被静默丢弃。那是把数据丢失伪装成去重。
处理器只做了一个SET和一个XADD。几毫秒的事,在重负载下也能轻松在5秒限制内。
第四步:幂等Worker
Webhook是至少一次投递,不是精确一次。这意味着你的Worker会看到重复和乱序事件,但必须产生正确的状态。这个属性叫幂等性,它让重试是安全的。
经典bug:进球#3的事件被重试,在"比赛结束"之后到达。简陋的Worker会覆盖最终状态,比赛看起来又变回直播了。
修复方案是序列门控:只有当事件的序列号高于该比赛上次应用的序列号时才应用。而且检查和写入必须是原子的,否则两个Worker会竞争。用Redis Lua脚本来实现:
worker.py:
import json
import logging
import sys
import time import redis from config import DLQ, GROUP, REDIS_URL, STREAM logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
log = logging.getLogger("worker") r = redis.from_url(REDIS_URL, decode_responses=True) MAX_ATTEMPTS = 5 # One atomic script: gate on sequence, write state, update the live set, publish.
# Because it's a single script, state and side effects can never diverge,
# even if the worker crashes mid-way. APPLY_LUA = """
local current = tonumber(redis.call('HGET', KEYS[1], 'seq') or '-1')
local incoming = tonumber(ARGV[1])
if incoming <= current then return 0 end redis.call('HSET', KEYS[1], 'seq', incoming, unpack(ARGV, 5)) if ARGV[3] == 'finished' then redis.call('SREM', KEYS[2], ARGV[2]) redis.call('EXPIRE', KEYS[1], 21600)
elseif ARGV[3] == 'in_play' then redis.call('SADD', KEYS[2], ARGV[2])
end redis.call('PUBLISH', 'updates', ARGV[4])
return 1
"""
apply_event = r.register_script(APPLY_LUA) def to_fields(ev: dict) -> dict: kind = ev["event"] fields = { "match_id": ev["match_id"], "sport": ev.get("sport"), "last_event": kind, "updated_at": ev.get("timestamp"), } if kind == "match.started": fields["status"] = "in_play" elif kind == "match.goal": fields.update( status="in_play", last_goal_team=ev.get("team"), last_goal_minute=ev.get("minute"), score=ev.get("score"), ) elif kind == "match.finished": fields.update(status="finished", final_score=ev.get("final_score")) # Unknown event types still update last_event, so new types don't crash us. return {k: v for k, v in fields.items() if v is not None} def handle(ev: dict) -> bool: fields = to_fields(ev) flat = [x for pair in fields.items() for x in pair] payload = json.dumps({"match_id": ev["match_id"], **fields}) applied = apply_event( keys=[f"match:{ev['match_id']}", "live:ids"], args=[int(ev["sequence"]), ev["match_id"], fields.get("status", ""), payload, *flat], ) return bool(applied) # False = stale or duplicate, safely ignored def process(msg_id: str, data: dict) -> None: try: handle(json.loads(data["payload"])) r.xack(STREAM, GROUP, msg_id) except Exception: log.exception("failed processing %s", msg_id) # Leave it pending; reclaim() retries it. After MAX_ATTEMPTS, park it. pending = r.xpending_range(STREAM, GROUP, min=msg_id, max=msg_id, count=1) if pending and pending[0]["times_delivered"] >= MAX_ATTEMPTS: r.xadd(DLQ, data) r.xack(STREAM, GROUP, msg_id) log.error("moved %s to dead-letter stream", msg_id) def reclaim(consumer: str) -> None: """Pick up messages a crashed worker never acknowledged.""" start = "0-0" while True: start, claimed, *_ = r.xautoclaim( STREAM, GROUP, consumer, min_idle_time=30_000, start_id=start, count=100 ) for msg_id, data in claimed: process(msg_id, data) if start == "0-0": break def main(consumer: str) -> None: try: r.xgroup_create(STREAM, GROUP, id="0", mkstream=True) except redis.ResponseError as exc: if "BUSYGROUP" not in str(exc): raise log.info("worker %s started", consumer) last_reclaim = 0.0 while True: if time.time() - last_reclaim > 15: reclaim(consumer) last_reclaim = time.time() resp = r.xreadgroup(GROUP, consumer, {STREAM: ">"}, count=100, block=2000) for _, messages in resp or []: for msg_id, data in messages: process(msg_id, data) if __name__ == "__main__": main(sys.argv[1] if len(sys.argv) > 1 else "worker-1") 为什么这个设计站得住:
消费者组把流分给多个Worker,横向扩展就是起另一个不同名字的进程。
成功后ACK意味着Worker崩溃时消息保持pending,reclaim()会重试。不会丢失任何东西。
死信流捕获毒消息,所以一个坏事件不会永久阻塞队列。
Lua脚本把"检查序列然后写入"变成一个原子步骤,并发不会破坏比赛状态。
比赛日横向扩展Worker就这这么简单:
python worker.py w1 &
python worker.py w2 &
python worker.py w3 & 第五步:一个有速率限制、带缓存的REST客户端
即使Webhook是主路径,你还是需要调用REST API:查赛程、积分榜、球队数据,以及安全垫轮询器。这个客户端得是个守规矩的公民。
大多数限速失败来自四个错误:忽略响应头、立即重试、步调一致地重试、让多个进程各以为自己拥有整个预算。这四个全给它修好。
API参考文档记录了这些构建块:
每个响应都有X-RateLimit-Limit、X-RateLimit-Remaining和X-RateLimit-Reset
429带Retry-After头
可安全重试的500和503响应
403是"你的套餐不包括这个"(不要重试)
client.py:
import json
import logging
import random
import time import redis
import requests from config import API_KEY, BASE_URL, BUDGET_PER_MIN, REDIS_URL log = logging.getLogger("client")
r = redis.from_url(REDIS_URL, decode_responses=True) session = requests.Session()
session.headers.update({"Authorization": f"Bearer {API_KEY}"}) class ApiError(Exception): pass def backoff(attempt: int, cap: float = 30.0) -> float: """Exponential backoff with FULL jitter, so clients don't retry in lockstep.""" return random.uniform(0, min(cap, 2 ** attempt)) def acquire() -> float: """Shared budget across ALL worker processes. Returns seconds to wait (0 = go).""" now = time.time() # A global pause (set after a 429 or a nearly-empty budget) beats everything. pause_until = r.get("rl:pause_until") if pause_until and float(pause_until) > now: return float(pause_until) - now key = f"rl:{int(now // 60)}" used = r.incr(key) if used == 1: r.expire(key, 120) return 0.0 if used <= BUDGET_PER_MIN else 60 - (now % 60) def pause_everyone(seconds: float) -> None: until = time.time() + min(seconds, 60) r.set("rl:pause_until", until, ex=int(seconds) + 2) def api_get(path: str, params: dict | None = None, retries: int = 4) -> dict: url = f"{BASE_URL}/{path.lstrip('/')}" for attempt in range(retries + 1): while wait := acquire(): time.sleep(wait + random.uniform(0, 0.5)) try: resp = session.get(url, params=params, timeout=(3, 10)) except requests.RequestException as exc: log.warning("network error on %s: %s", path, exc) time.sleep(backoff(attempt)) continue # Back off *before* we hit the wall, for every worker at once. remaining = resp.headers.get("X-RateLimit-Remaining") reset = resp.headers.get("X-RateLimit-Reset") if remaining is not None and reset and int(remaining) <= 2: pause_everyone(max(0, float(reset) - time.time())) if resp.status_code == 429: delay = float(resp.headers.get("Retry-After", backoff(attempt, cap=60))) log.warning("429 on %s, pausing all workers for %.1fs", path, delay) pause_everyone(delay) time.sleep(delay + random.uniform(0, 1)) continue if resp.status_code >= 500: time.sleep(backoff(attempt)) continue if resp.status_code in (400, 401, 403, 404): # Retrying won't fix a bad key, a plan limit, or a typo. raise ApiError(f"{resp.status_code} on {path}: {resp.text[:200]}") resp.raise_for_status() return resp.json() raise ApiError(f"gave up on {path} after {retries + 1} attempts") 两个关键思想:
pause_everyone。当一个Worker收到429,所有Worker都停。否则十个Worker各自独立发现限制,各自重试冲进去。
全抖动。random.uniform(0, 2**attempt)把重试分散开。普通指数退避让每个客户端在同一时刻重试,又重建了你在努力避免的流量峰值。
旁路缓存、单飞、过期重刷
现在这个部分保护你免受你自己的用户伤害。把它加到client.py:
def cached_get(path: str, params: dict | None = None, ttl: int = 10, stale_ttl: int = 600): """ - Fresh hit: return immediately (no upstream call). - Miss: exactly ONE caller fetches; others get stale data or wait briefly. - Upstream failing: serve stale data instead of an error. """ key = f"http:{path}:{json.dumps(params or {}, sort_keys=True)}" fresh = r.get(key) if fresh: r.incr("stats:cache_hit") return json.loads(fresh) r.incr("stats:cache_miss") owns_lock = bool(r.set(f"lock:{key}", 1, nx=True, ex=10)) if not owns_lock: stale = r.get(f"{key}:stale") if stale: return json.loads(stale) # serve old data, don't pile on for _ in range(20): # wait up to ~2s for the leader time.sleep(0.1) fresh = r.get(key) if fresh: return json.loads(fresh) # Leader never delivered; fall through and fetch ourselves. try: body = api_get(path, params) payload = json.dumps(body) pipe = r.pipeline() pipe.set(key, payload, ex=ttl) pipe.set(f"{key}:stale", payload, ex=ttl + stale_ttl) pipe.execute() return body except Exception: stale = r.get(f"{key}:stale") if stale: log.warning("upstream failed for %s, serving stale", path) return json.loads(stale) raise finally: if owns_lock: r.delete(f"lock:{key}") 每个部分带来什么:
单飞(那个锁):如果5000个请求在同一毫秒内缓存未命中,只有一个去上游。这防止了"缓存雪崩"——缓存过期时把系统打垮。
过期重刷:调用方立即拿到稍旧的数据,同时一个调用方去刷新它。对体育比分来说,"旧8秒"比"报错"好。
故障时返回过期数据:如果上游在限速或宕机,你继续提供最后一次正确的答案。
根据数据实际变化速度选择TTL
API参考文档标记赛程为可缓存,直播端点为不可缓存。用它来设置你自己的策略:
数据 建议TTL 过期窗口 原因
赛程/日程 5分钟 1小时 发布后几乎不变
积分榜 60秒(比赛日) 10分钟 只有结果确认后才变化
直播比赛(批量) 8秒 2分钟 变化快,但一次调用服务所有人
球队、联赛 24小时 24小时 实际上静态
这些是起点。测量你的缓存命中率(后面会暴露),然后调整。
第六步:对账轮询器(你的安全垫)
Webhook很棒,但没有推送系统是完美的。你的端点可能在部署时宕机。网络抖动可能吞掉一次交付。你这边的bug可能丢掉一个事件。
对账轮询器悄悄比较实际情况和你的状态,修复缺口。它跑得慢因为它是安全垫,不是主路径。
poller.py:
import logging
import time
from datetime import datetime from client import cached_get, r
from config import LIVE_PATH logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
log = logging.getLogger("poller") # The 13 sports. Confirm exact slugs in the documentation.
SPORTS = [ "football", "basketball", "american-football", "cricket", "tennis", "baseball", "esports", "combat-sports", "volleyball", "handball", "ice-hockey", "golf", "horse-racing",
] def parse_ts(value): return datetime.fromisoformat(value.replace("Z", "+00:00")) if value else None def reconcile(sport: str) -> int: body = cached_get(LIVE_PATH.format(sport=sport), ttl=8, stale_ttl=120) items = body.get("data", body) if isinstance(body, dict) else body if isinstance(items, dict): # some responses return a single object items = [items] repaired = 0 for m in items or []: match_id = m.get("fixture_id") or m.get("match_id") if not match_id: continue key = f"match:{match_id}" ours = r.hgetall(key) api_ts, our_ts = parse_ts(m.get("updated_at")), parse_ts(ours.get("updated_at")) missing = not ours behind = bool(api_ts and our_ts and api_ts > our_ts) if not (missing or behind): continue score = m.get("score") or {} fields = { "match_id": match_id, "sport": sport, "status": m.get("status", "in_play"), "minute": m.get("minute"), "score": f"{score.get('home')}-{score.get('away')}" if score else None, "updated_at": m.get("updated_at"), "source": "poll", } # We deliberately never touch 'seq', so a later webhook still wins. r.hset(key, mapping={k: v for k, v in fields.items() if v is not None}) r.sadd("live:ids", match_id) repaired += 1 return repaired def main(interval: int = 30) -> None: while True: started = time.time() for sport in SPORTS: try: n = reconcile(sport) if n: log.warning("repaired %d %s matches (webhooks missed something)", n, sport) except Exception: log.exception("reconcile failed for %s", sport) time.sleep(max(0, interval - (time.time() - started))) if __name__ == "__main__": main() 注意那行日志。如果轮询器经常在修复东西,说明你的Webhook路径有问题。轮询器既是安全垫,也是早期预警系统。
后面值得做的一个改进:跳过今天没有比赛的运动项目(查一次赛程,缓存五分钟),这样安静的运动项目不消耗任何请求。
还要记住,轮询器和所有其他东西共享同一个速率预算,通过acquire()。它不能饿死你的其他系统。
第七步:让你的用户只从缓存读数据
现在到了收获的时候。给app.py加两个只从Redis读的端点:
@app.get("/matches/live")
async def live_matches(): ids = await r.smembers("live:ids") pipe = r.pipeline() for match_id in ids: pipe.hgetall(f"match:{match_id}") rows = await pipe.execute() return [row for row in rows if row] @app.get("/stream")
async def stream(): """Server-Sent Events: one Redis subscription per client, zero upstream calls.""" async def events(): pubsub = r.pubsub() await pubsub.subscribe("updates") try: async for msg in pubsub.listen(): if msg["type"] == "message": yield f"data: {msg['data']}\n\n" finally: await pubsub.unsubscribe("updates") await pubsub.aclose() return StreamingResponse(events(), media_type="text/event-stream") @app.get("/healthz")
async def healthz(): try: groups = await r.xinfo_groups(STREAM) group = next((g for g in groups if g["name"] == GROUP), {}) except Exception: group = {} return { "stream_length": await r.xlen(STREAM), "pending": group.get("pending"), "lag": group.get("lag"), "dead_letters": await r.xlen(DLQ), "live_matches": await r.scard("live:ids"), } 浏览器三行代码就能消费这个流:
const es = new EventSource("/stream");
es.onmessage = (e) => updateScoreboard(JSON.parse(e.data)); 不管你有50个用户还是50000个用户,你对上游API的调用量完全相同。横向扩展用户现在只花Redis读操作的钱,Redis读很便宜而且可以水平扩展。
启动服务:
uvicorn app:app --workers 2 --port 8000 懒得自己搭广播?
如果你只需要展示比分,不需要自定义逻辑,可以跳过整个这层。平台的小组件嵌入一个实时比分牌、比赛中心或赔率板,用一个script标签搞定,背后的数据源是同一个实时流。只有当你的产品需要组件给不了的数据、界面或逻辑控制权时,自己搭广播基础设施才值得。
第八步:注册Webhook(并在宕机后重放)
你的端点必须可以通过HTTPS访问。开发期间,用任意隧道工具把localhost:8000暴露出来。然后注册:
register.py:
import requests
from config import API_KEY, BASE_URL, WEBHOOK_SECRET resp = requests.post( f"{BASE_URL}/webhooks", headers={"Authorization": f"Bearer {API_KEY}"}, json={ "url": "https://your-domain.example/webhook", "events": ["match.started", "match.goal", "match.finished"], "secret": WEBHOOK_SECRET, # optional: auto-generated if omitted # "sport": "football", # optional: limit to one sport }, timeout=15,
)
resp.raise_for_status()
print(resp.json()) 只订阅你实际用到的事件类型。每个不用的事件类型都是你花钱接收和处理的无用流量。
当你宕机时会发生什么?
文档描述了自动重试计划:一次立即尝试,然后大约30秒、2分钟、10分钟、1小时,之后交付标记为失败。任何非2xx响应或超时都会触发重试。
这覆盖了短暂的抖动。长时间宕机不会丢失任何东西:遗漏的交付可以被列出并重放。
def replay_missed(webhook_id: str) -> None: headers = {"Authorization": f"Bearer {API_KEY}"} base = f"{BASE_URL}/webhooks/{webhook_id}" print(requests.get(f"{base}/deliveries", headers=headers, timeout=15).json()) requests.post(f"{base}/replay", headers=headers, timeout=15).raise_for_status() 查看文档了解任何过滤参数。重点是:只有你的Worker是幂等的,重放才是安全的。重放的事件携带相同的event_id和sequence,所以已经应用过的事件会被简单忽略。这就是第四步的回报。
一个好的宕机运行手册:
修复并重新部署你的接收器。
触发重放,覆盖你遗漏的时间窗口。
看着轮询器的"修复了N场比赛"日志行恢复到零。
第九步:WebSocket替代方案
如果你更喜欢流式,这个WebSocket API是摄入步骤的替代方案。下游所有东西(队列、Worker、缓存、读API)完全不变。把队列放在摄入和处理之间,这就是好处。
文档显示这些行为:
连接时认证,然后发送{"action": "subscribe", "channel": "football.live"}
每个订阅可选league和match_id过滤器
通过ping/pong帧心跳
关闭码:1000正常、4001密钥无效或缺失、4008超过速率限制、1006异常关闭
一个坑:文档中的WebSocket消息没有event_id或sequence。我们自己推导两者,这样同一个Worker可以处理它:
ws_consumer.py:
import asyncio
import hashlib
import json
import logging
import random
from datetime import datetime import redis.asyncio as aioredis
import websockets from config import API_KEY, REDIS_URL, STREAM log = logging.getLogger("ws")
r = aioredis.from_url(REDIS_URL, decode_responses=True) # Confirm the exact auth parameter (api_key vs token) in the docs/sandbox.
URL = f"wss://stream.orbistats.com/v1?api_key={API_KEY}"
CHANNELS = ["football.live", "cricket.live"] def normalise(msg: dict) -> dict: """Give a WebSocket message the same shape the worker expects.""" ts = msg["timestamp"] fingerprint = f'{msg["match_id"]}|{msg["event"]}|{msg.get("minute")}|{msg.get("team")}|{ts}' dt = datetime.fromisoformat(ts.replace("Z", "+00:00")) return { **msg, "event