DEV Community

Cover image for Scaling a Sports Data Consumer for Match-Day Spikes: Queues, Caching and Rate Limits
orbistats
orbistats

Posted on

Scaling a Sports Data Consumer for Match-Day Spikes: Queues, Caching and Rate Limits

Your sports app works perfectly on a quiet Tuesday.

Then Saturday, 3:00 PM arrives. Dozens of matches kick off within the same minute. Goals land in parallel across several leagues. Ten minutes later, your users are all refreshing at once, and your single-threaded poller is hammering an API that just returned its first 429.

Everything that was "fine in testing" turns out to depend on traffic being polite.

This tutorial builds a sports data consumer that stays up when traffic isn't polite. We'll build it in Python with Redis, and every piece is something you can reuse for any real-time feed, not just sports.

By the end you'll have:

A webhook receiver that verifies signatures and acknowledges in milliseconds
A durable queue (Redis Streams) that absorbs bursts instead of dropping them
Idempotent workers that survive duplicates, retries and out-of-order events
A rate-limited API client with shared budgets, backoff, jitter and Retry-After handling
Cache-aside with single-flight and stale-while-revalidate, so a thousand users cost one API call
A reconciliation poller that heals anything a webhook missed
A spike simulator that proves all of the above actually works

Let's build it.

Why Match Days Break Naive Consumers

Sports traffic isn't smooth. It's synchronized, and synchronization is what kills systems.

Kickoffs cluster. A Premier League Saturday puts a whole slate of matches on the same clock. Fixtures that were quiet for a week become live all at once.
Events cluster. A burst of goals, cards and substitutions across many matches arrives in the same few seconds, then nothing for a minute.
Final whistles cluster. Matches that kicked off together end together, so a wave of match.finished events lands at nearly the same instant.
Your users cluster too. A goal in a big match sends tens of thousands of people to refresh simultaneously.
Sports overlap. A cricket T20 final, a tennis tournament day, and a football slate can all peak together, so the spike isn't confined to one sport.

A naive consumer fails in predictable ways:

Naive approach What breaks on match day
One API call per user request Your users become your rate-limit problem
Poll every match every few seconds Request count scales with matches, not with value
Process events inline in the HTTP handler A slow database write times out the sender
Trust arrival order A retried "goal" lands after the "final score" and flips the match back to live
Retry immediately on 429 You turn a rate limit into a thundering herd
No dedupe Retries double-count goals

Each of those has a boring, well-understood fix. We'll apply them one at a time.

The Architecture
text
┌────────────────────────────┐
│ Orbistats feed │
└──────┬───────────────┬─────┘
webhooks (push)│ │ REST (pull, safety net)
▼ ▼
┌───────────────┐ ┌───────────────┐
│ FastAPI │ │ Reconciliation│
│ receiver │ │ poller │
│ verify+dedupe │ └───────┬───────┘
└───────┬───────┘ │
▼ │
┌───────────────┐ │
│ Redis Stream │ │
│ (durable queue│ │
└───────┬───────┘ │
▼ ▼
┌───────────────────────────────┐
│ Workers (N) → atomic Lua │
│ state: match:{id} + live:ids │──► pub/sub "updates"
└───────────────┬───────────────┘
▼
┌───────────────────────────────┐
│ Read API + SSE (your users) │
│ serves from cache, never from │
│ the upstream API │
└───────────────────────────────┘

The golden rule behind the whole design: your users should never cause a call to the upstream API. Users read from your cache. The upstream API is fed by your workers, at a rate you control.

Choosing How to Receive Data

Before writing code, decide how data gets into your system. The Live Scores API delivers the same events three ways:

Method Best for Trade-off
REST polling Simple jobs, reconciliation Cost scales with poll frequency; always a little late
WebSocket Latency-sensitive live UIs You own reconnect logic and connection limits
Webhooks Event-driven backends You need a public HTTPS endpoint that answers quickly

For match-day scale I recommend webhooks as the primary path with REST as a safety net:

Webhooks are push-based, so your request count doesn't grow with the number of live matches.
They come with the properties a robust consumer needs: signed payloads, automatic retries with backoff, an event_id for dedupe, and a per-match sequence for ordering.
If your endpoint is down for a while, deliveries can be listed and replayed afterwards.

We'll also show the WebSocket variant, since many teams prefer it.

Capacity Math Before You Write Code

Do this on paper first. These numbers are illustrative assumptions, not Orbistats limits:

text
Naive: poll each live match every 10 s
300 live matches ÷ 10 s = 30 requests/second = 1,800 requests/minute

Better: one bulk "live" call per sport every 10 s
13 sports ÷ 10 s ≈ 78 requests/minute

Best: webhooks as primary, bulk poll every 30 s as a safety net
13 sports ÷ 30 s ≈ 26 requests/minute

Your users:
50,000 users refreshing every 10 s = 5,000 reads/second
→ that must hit YOUR cache, not the upstream API

The same data costs 1,800 requests per minute done badly and 26 done well. Plan limits differ by tier, so check the pricing page and the rate-limit headers on real responses, then size your BUDGET_PER_MIN accordingly.

Step 1: Setup
Request an API key on the sign-up page.
Read the documentation for auth and conventions, then try the endpoints in the sandbox so you see real payloads before coding.
Keep the API reference open. It documents the X-RateLimit-* headers, the 429 and Retry-After behaviour, and the standard response envelope (data, meta, errors).

Project layout:

text
matchday/
├── .env
├── docker-compose.yml
├── config.py
├── app.py # webhook receiver + read API + SSE
├── worker.py # stream consumer, idempotent state updates
├── client.py # rate-limited, cached REST client
├── poller.py # reconciliation safety net
├── ws_consumer.py # optional WebSocket ingestion
├── register.py # register the webhook
└── simulate.py # spike simulator + verifier

Install:

bash
python -m venv .venv && source .venv/bin/activate
pip install fastapi uvicorn "redis>=5" requests httpx websockets python-dotenv

docker-compose.yml (Redis 7 is needed for the stream commands we use):

yaml
services:
redis:
image: redis:7-alpine
command: ["redis-server", "--appendonly", "yes"]
ports: ["6379:6379"]
volumes: ["redis-data:/data"]

volumes:
redis-data:

.env:

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

Generate a strong webhook secret:

bash
python -c "import secrets; print(secrets.token_urlsafe(32))"

Start Redis:

bash
docker compose up -d
Step 2: Shared Configuration

config.py:

python
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
Step 3: The Webhook Receiver (Acknowledge Fast, Do Nothing Else)

The webhook docs set the rule: respond within 5 seconds or the delivery counts as failed and is retried. The lesson is simple: do the minimum in the HTTP handler. Verify, dedupe, enqueue, return 200. Everything slow happens later, in a worker.

Each delivery is a signed POST. The signature is an HMAC-SHA256 of the raw request body using your secret, sent in the X-Orbistats-Signature header. Here's the shape of a payload from the docs:

json
{
"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"
}

Two fields matter most for scaling:

event_id is stable across retries, so dedupe on it.
sequence is monotonic per match_id, so order by it, not by arrival time.

app.py:

python
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
Enter fullscreen mode Exit fullscreen mode

Three details worth stealing:

We sign-check before parsing. Never trust a payload you haven't authenticated.
Duplicates get a 200. The sender should stop retrying, even though we ignored the event.
The "forget on failure" branch. Without it, a failed enqueue would mark the event as seen and the retry would be silently discarded. That's data loss disguised as dedupe.

The handler does one SET and one XADD. That's a few milliseconds, comfortably inside the 5-second limit even under heavy load.

Step 4: Idempotent Workers

Webhooks are at-least-once, not exactly-once. That means your worker will see duplicates and out-of-order events, and it must produce the correct state anyway. This property is called idempotency, and it's what makes retries safe.

The classic bug: the event for "goal #3" is retried and arrives after "match finished". A naive worker overwrites the final state and the match looks live again.

The fix is sequence gating: only apply an event if its sequence is higher than the last one applied for that match. And the check-and-write has to be atomic, otherwise two workers can race. A Redis Lua script gives us that:

worker.py:

python
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)
Enter fullscreen mode Exit fullscreen mode

if name == "main":
main(sys.argv[1] if len(sys.argv) > 1 else "worker-1")

Why this design holds up:

Consumer groups split the stream across workers, so scaling out is just launching another process with a different name.
Ack-after-success means a worker crash leaves the message pending, and reclaim() retries it later. Nothing is lost.
A dead-letter stream catches poison messages, so one bad event can't block the queue forever.
The Lua script makes "check sequence, then write" a single atomic step, so concurrency can't corrupt a match.

Scaling workers on match day is literally:

bash
python worker.py w1 &
python worker.py w2 &
python worker.py w3 &
Step 5: A Rate-Limit-Aware, Cached REST Client

Even with webhooks as the primary path, you'll still call the REST API: for fixtures, standings, team data, and the safety-net poller. That client needs to be a good citizen.

Most rate-limit failures come from four mistakes: ignoring the headers, retrying instantly, retrying in lockstep, and letting many processes each think they own the whole budget. We'll fix all four.

The API reference documents these building blocks:

X-RateLimit-Limit, X-RateLimit-Remaining and X-RateLimit-Reset on every response
429 with a Retry-After header
Safe-to-retry 500 and 503 responses
403 for "your plan doesn't include this" (do not retry that)

client.py:

python
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)
Enter fullscreen mode Exit fullscreen mode

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")
Enter fullscreen mode Exit fullscreen mode

The two ideas that matter most:

pause_everyone. When one worker sees a 429, all workers stop. Otherwise ten workers each discover the limit independently and each one retries into it.
Full jitter. random.uniform(0, 2**attempt) spreads retries out. Plain exponential backoff makes every client retry at the same instant, which recreates the spike you were avoiding.
Cache-Aside, Single-Flight, Stale-While-Revalidate

Now the part that protects you from your own users. Add this to client.py:

python
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}")
Enter fullscreen mode Exit fullscreen mode

What each part buys you:

Single-flight (the lock): if 5,000 requests miss the cache in the same millisecond, one goes upstream. This prevents the "cache stampede" that takes down systems right when the cache expires.
Stale-while-revalidate: callers get slightly old data instantly while one caller refreshes it. For sports scores, "8 seconds old" is better than "an error".
Stale-on-error: if the upstream is rate-limiting or down, you keep serving the last good answer.
Pick TTLs by how fast the data actually changes

The API reference marks fixtures as cacheable and the live endpoint as not cached. Use that to set your own policy:

Data Suggested TTL Stale window Why
Fixtures / schedules 5 min 1 hour Rarely change once published
Standings 60 s (on match day) 10 min Change only after a result is confirmed
Live matches (bulk) 8 s 2 min Fast-moving, but one call serves everyone
Teams, competitions 24 h 24 h Effectively static

These are starting points. Measure your cache hit ratio (we expose it later) and tune.

Step 6: The Reconciliation Poller (Your Safety Net)

Webhooks are excellent, but no push system is perfect. Your endpoint might be down during a deployment. A network blip might swallow a delivery. A bug on your side might drop an event.

A reconciliation poller quietly compares reality with your state and repairs gaps. It runs slowly because it's a safety net, not the primary path.

poller.py:

python
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
Enter fullscreen mode Exit fullscreen mode

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()

Notice the log line. If the poller repairs things often, your webhook path has a problem. The poller is both a safety net and an early-warning system.

One improvement worth making later: skip sports with no matches today (check fixtures once, cached for five minutes), so quiet sports cost you zero requests.

Also remember that the poller shares the same rate budget as everything else, via acquire(). It can't starve the rest of your system.

Step 7: Serve Your Own Users From Cache

Now the payoff. Add two endpoints to app.py that read only from Redis:

python
@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")
Enter fullscreen mode Exit fullscreen mode

@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"),
}

Browsers can consume the stream with three lines:

javascript
const es = new EventSource("/stream");
es.onmessage = (e) => updateScoreboard(JSON.parse(e.data));

Whether you have 50 users or 50,000, your upstream API usage is identical. Scaling users now only costs Redis reads, which are cheap and scale horizontally.

Run it:

bash
uvicorn app:app --workers 2 --port 8000
Can't be bothered to build the fan-out yourself?

If you only need to display scores and don't need custom logic, you can skip this whole layer. The platform's Widgets embed a live scoreboard, match center or odds board with a script tag, powered by the same live feed. Building fan-out infrastructure is only worth it when your product needs control over data, UI or logic that a widget can't give you.

Step 8: Register the Webhook (and Replay After Outages)

Your endpoint must be reachable over HTTPS. During development, use any tunnelling tool to expose localhost:8000. Then register it:

register.py:

python
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())

Subscribe only to the events you actually use. Every unused event type is traffic you pay to receive and process.

What happens when you're down?

The docs describe an automatic retry schedule: an immediate attempt, then roughly 30 seconds, 2 minutes, 10 minutes, and 1 hour, after which the delivery is marked failed. Any non-2xx response or a timeout triggers a retry.

That covers short blips. For longer outages, nothing is lost: missed deliveries can be listed and replayed.

python
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()
Enter fullscreen mode Exit fullscreen mode

Check the docs for any filtering parameters. The important point: replay is only safe because your workers are idempotent. Replayed events carry the same event_id and sequence, so already-applied events are simply ignored. That's the payoff for Step 4.

A good outage runbook:

Fix and redeploy your receiver.
Trigger a replay for the window you missed.
Watch the poller's "repaired N matches" log line fall back to zero.
Step 9: The WebSocket Alternative

If you prefer streaming, the WebSocket API is a drop-in replacement for the ingestion step. Everything downstream (queue, workers, cache, read API) stays exactly the same. That's the benefit of putting a queue between ingestion and processing.

The docs show these behaviours:

Authenticate on connect, then send {"action": "subscribe", "channel": "football.live"}
Optional league and match_id filters per subscription
Heartbeats via ping/pong frames
Close codes: 1000 normal, 4001 invalid or missing key, 4008 rate limit exceeded, 1006 abnormal closure

One catch: the documented WebSocket message has no event_id or sequence. We derive both so the same worker can process it:

ws_consumer.py:

python
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_id": "ws_" + hashlib.sha1(fingerprint.encode()).hexdigest()[:16],
"sequence": int(dt.timestamp() * 1000), # ordering by event time
}

async def run() -> None:
attempt = 0
while True:
forced_delay = None
try:
async with websockets.connect(URL, ping_interval=20, ping_timeout=20) as ws:
attempt = 0
for channel in CHANNELS:
await ws.send(json.dumps({"action": "subscribe", "channel": channel}))
log.info("connected, subscribed to %s", CHANNELS)

            async for raw in ws:
                msg = json.loads(raw)
                if "match_id" not in msg or "event" not in msg:
                    continue                    # skip acks and housekeeping frames
                await r.xadd(STREAM, {"payload": json.dumps(normalise(msg))},
                             maxlen=500_000, approximate=True)

    except websockets.ConnectionClosed as exc:
        code = exc.rcvd.code if exc.rcvd else 1006
        if code == 4001:
            raise SystemExit("API key rejected (close code 4001)")
        forced_delay = 60 if code == 4008 else None   # rate limited: back off hard
        log.warning("connection closed (%s)", code)
    except OSError as exc:
        log.warning("network error: %s", exc)

    attempt += 1
    await asyncio.sleep(forced_delay or (min(30, 2 ** attempt) + random.random()))
Enter fullscreen mode Exit fullscreen mode

if name == "main":
asyncio.run(run())

Two cautions:

Don't use webhooks and WebSockets as co-equal primary sources for the same match. Their ordering fields come from different schemes (a server counter versus a timestamp), so mixing them breaks sequence gating. Pick one primary source and keep REST as the safety net.
This path skips the receiver's dedupe. That's fine here, because the worker's sequence gate already ignores duplicates (same timestamp means same sequence, which is not greater than the stored one).

Also read the close-code handling carefully. 4001 is fatal (retrying will never help), 4008 means you're being rate limited (back off for a long time), and 1006 is an ordinary drop (reconnect with jittered exponential backoff).

Step 10: Prove It With a Spike Simulator

Everything above is theory until you test it. This simulator generates a worst-case match day:

300 matches with 6 events each
20% duplicate deliveries (retries)
Fully shuffled arrival order (so "finished" often arrives before "goal")
Hundreds of concurrent, properly signed requests

simulate.py:

python
import asyncio
import hashlib
import hmac
import json
import random
import time

import httpx
import redis

from config import REDIS_URL, STREAM, WEBHOOK_SECRET

URL = "http://localhost:8000/webhook"
MATCHES, PER_MATCH = 300, 6
r = redis.from_url(REDIS_URL, decode_responses=True)

def sign(raw: bytes) -> str:
return "sha256=" + hmac.new(WEBHOOK_SECRET.encode(), raw, hashlib.sha256).hexdigest()

def build_batch() -> list[dict]:
events = []
for m in range(MATCHES):
match_id = 60000 + m
for seq in range(1, PER_MATCH + 1):
last = seq == PER_MATCH
events.append({
"event": "match.finished" if last else "match.goal",
"event_id": f"evt_{match_id}{seq}",
"delivery_id": f"dlv
{match_id}_{seq}",
"sequence": seq,
"match_id": match_id,
"sport": "football",
**({"final_score": "3-2"} if last else {"team": "Home", "minute": seq * 10}),
"timestamp": "2026-09-14T16:52:11Z",
})
retries = random.sample(events, len(events) // 5) # duplicate deliveries
batch = events + retries
random.shuffle(batch) # out-of-order arrival
return batch

async def blast(batch: list[dict], concurrency: int = 100) -> None:
sem = asyncio.Semaphore(concurrency)

async with httpx.AsyncClient(timeout=5) as client:
    async def send(ev: dict) -> int:
        raw = json.dumps(ev).encode()
        async with sem:
            resp = await client.post(
                URL, content=raw,
                headers={"X-Orbistats-Signature": sign(raw),
                         "Content-Type": "application/json"},
            )
        return resp.status_code

    t0 = time.time()
    codes = await asyncio.gather(*(send(e) for e in batch))
    elapsed = time.time() - t0

ok = sum(c == 200 for c in codes)
print(f"sent {len(batch)} deliveries in {elapsed:.1f}s "
      f"({len(batch) / elapsed:.0f}/s), {ok} acknowledged with 200")
Enter fullscreen mode Exit fullscreen mode

def verify(timeout: int = 30) -> None:
deadline = time.time() + timeout
while time.time() < deadline:
good = sum(
1 for m in range(MATCHES)
if (h := r.hgetall(f"match:{60000 + m}")).get("status") == "finished"
and int(h.get("seq", 0)) == PER_MATCH
)
if good == MATCHES:
break
time.sleep(1)

print(f"matches in correct final state: {good}/{MATCHES}")
print(f"unique events enqueued: {r.xlen(STREAM)} (expected {MATCHES * PER_MATCH})")
print(f"live set size: {r.scard('live:ids')} (expected 0)")
print(f"dead letters: {r.xlen('events:dead')} (expected 0)")
Enter fullscreen mode Exit fullscreen mode

if name == "main":
asyncio.run(blast(build_batch()))
verify()

Run it (with Redis, the receiver, and at least one worker already running):

bash
python simulate.py

What a healthy run looks like:

All 300 matches end in the correct final state, even though events arrived shuffled
Unique events enqueued equals 1,800, which proves the 20% duplicates were dropped at the door
The live set is empty, because every match ended
Dead letters is zero

If you want to re-run, clear state first so the "seen" keys don't suppress your events:

bash
redis-cli FLUSHDB

Try breaking things on purpose. Kill a worker mid-run and watch reclaim() finish its work. Stop all workers, send the batch, then start them and watch the backlog drain. That's the queue doing its job.

Observability: What to Watch on Match Day

You can't fix what you can't see, and on match day you won't have time to dig. Watch these:

Signal Where it comes from Why it matters
Queue lag / pending /healthz (lag, pending) Rising lag means workers can't keep up; add workers
Dead letters /healthz (dead_letters) Anything above zero is a bug to investigate
Webhook ack time (p95) Receiver logs / APM Must stay well under the 5-second limit
429 count Log line from api_get Your budget is too high or something is polling wildly
Cache hit ratio stats:cache_hit vs stats:cache_miss Low ratio means TTLs too short or keys too varied
Poller repairs Poller warning log Frequent repairs mean webhooks are being lost
Redis memory Redis INFO memory The stream and caches grow during peaks

A tiny cache-ratio helper:

python
hits = int(r.get("stats:cache_hit") or 0)
misses = int(r.get("stats:cache_miss") or 0)
print(f"cache hit ratio: {hits / max(1, hits + misses):.1%}")

Sensible starting alerts: consumer lag staying above a few hundred for more than a minute, any dead letters, and any sustained run of 429 responses. When your own error rates climb unexpectedly, check the Orbistats status page before you start debugging your own stack. Knowing it's upstream saves a lot of wasted panic.

Failure Modes and What Actually Happens
Failure Behaviour in this design
Worker crashes mid-event Message stays pending; reclaim() retries it; the atomic script prevents half-applied state
Duplicate delivery (retry) Dropped at the receiver by event_id, or ignored by the sequence gate
Out-of-order delivery Sequence gate keeps the newest state
Upstream returns 429 All workers pause; Retry-After honoured; stale cache served meanwhile
Upstream returns 5xx Jittered exponential backoff; stale cache served
Your endpoint down for hours Retries exhaust; use replay; poller repairs the gap
Redis restarts AOF persistence preserves the stream; consumers resume from their group position
Poison event Retried 5 times, then moved to the dead-letter stream
Sudden user surge Reads hit Redis only; upstream load is unchanged
Bad signature Rejected with 401, never enqueued
Match-Day Checklist

Run through this the day before a big fixture list:

Budget set below your real plan limit, using the headers you observe
Webhook subscribed only to events you use
Receiver deployed behind HTTPS with more than one process
At least two workers running, with a way to add more quickly
Poller running on a relaxed interval
Dead-letter alert configured
Redis persistence on and memory headroom confirmed
Spike simulator run against staging in the last week
Replay procedure written down, so nobody improvises during an outage
Stream maxlen large enough that a long outage can't trim unprocessed events
Secrets (API key, webhook secret) kept out of git
Ideas to Extend This Project

Once the foundation holds, there's a lot you can build on top of it:

Odds-movement pipeline. Subscribe to odds.moved and feed it through the same queue and workers. The Odds API uses the same delivery model, so the architecture carries over almost unchanged.
Low-latency pricing tools. If sub-second reaction matters, the pattern here is the same one Trading Desks rely on: a persistent stream, an atomic state layer, and no slow work in the hot path.
Live match centers for newsrooms. Publishers get the most from the cache-first read API and SSE stream. The Media & Publishers page describes that use case.
Per-sport tuning. Different sports spike differently. Cricket has long quiet stretches and sudden bursts, tennis has many parallel matches, and football clusters at fixed kickoff times. Give each its own TTLs and worker pool.
Backpressure and load shedding. When lag grows past a threshold, drop low-value events (like minute ticks) and keep goals and final results.
Autoscaling workers based on consumer lag rather than CPU.
Multi-region read replicas for the read API if your audience is global.
Wrapping Up

We took a consumer that works on a quiet Tuesday and made it survive a synchronized Saturday:

Acknowledge fast. The receiver verifies, dedupes and enqueues. Nothing slow happens inline.
Queue everything. Redis Streams absorb bursts, and consumer groups scale workers horizontally.
Make processing idempotent. Sequence gating and an atomic script make duplicates and reordering harmless.
Share the rate budget. One pause signal across all workers, jittered backoff, and Retry-After respected.
Cache like you mean it. Single-flight, stale-while-revalidate, and stale-on-error mean users never cost an upstream call.
Keep a safety net. A slow reconciliation poller and webhook replay mean "missed" doesn't mean "lost".
Test the ugly case. The simulator proves it with duplicates, shuffled order and concurrency.

If you remember one principle, make it this: decouple the rate at which data arrives from the rate at which you process it, and the rate at which users read from the rate at which you fetch. Queues and caches are just that idea applied twice.

If you build on this (an autoscaler, a metrics dashboard, a multi-sport version), share it in the comments. And if the sandbox returns a payload shape that differs from the examples here, paste it below and I'll help adapt to_fields().

Happy scaling! ⚡

Top comments (0)