sifting/io
Developer Tutorials
16 min readSiftingIO Team

WebSocket fan-out with Redis pub/sub: one market data connection for every service

Share one market data WebSocket across services: read max_conn and max_subs from the ack, republish ticks on Redis, snapshot for late joiners. Python code.

WebSocket fan-out with Redis pub/sub: one market data connection for every service

WebSocket fan-out with Redis pub/sub fixes a specific failure. An alerting service, a charting service and a pricing cache each open their own stream to wss://stream.sifting.io/ws/v1, and the second, fourth or eleventh one is refused with {"f":"error","code":"max_connections"}. Connection and subscription caps on the stream belong to the API key's tier, so every consumer that opens its own socket spends a budget that belongs to the whole system. The answer is one process that holds the upstream connection, subscribes to the union of what every service wants, and republishes each tick on an internal Redis channel. Services subscribe to Redis instead of to the vendor.

This post builds that process in Python. The hub authenticates once, reads max_conn and max_subs from the ack frame instead of hardcoding a tier, keeps the union of symbols subscribed, refuses a request locally when the union would exceed the budget, republishes every tick on md:tick:<product>:<symbol>, and stores the latest frame per symbol so a service that starts late gets a price immediately. A reconnect inside the hub is invisible to consumers apart from a pause. Frame handling is a pure function, so the error frames get a pytest file rather than a hope. If you have not used the stream before, Real-time FX and crypto quotes: REST snapshots and WebSocket streams covers subscribing to a single symbol; this post starts where that one ends.

The caps, and where the server tells you what they are#

The pricing page (checked 2026-09-29) lists a fixed number of concurrent WebSocket connections and symbol subscriptions per tier, and states that the limits apply per market:

TierConnectionsSymbol subscriptions
Free15
Builder3100
Pro101,000
Ultra50Unlimited
EnterpriseCustomUnlimited

Those numbers change, so the hub should not carry them in a config file. After the connection authenticates, the server sends an ack that carries them (example from the docs page):

{"f":"ack","op":"auth","tier":"pro","max_conn":10,"max_subs":1000,"active_conn":1}

max_subs is the budget for the connection the hub holds. active_conn is the number of connections the key has open right now, and it's the figure to log when a deploy starts getting refused: if it reads 3 on a Builder key while the hub is the only thing that should be connected, something else still is.

Two error frames matter for a fan-out. max_connections, message "Tier connection cap reached", is what a new connection gets when the tier's connection cap is already in use. max_subscriptions, message "Tier subscription cap reached", is what a subscribe gets when it would push the connection past max_subs. Both carry a limit field. What the docs do not say is what happens to the connections already open when a new one is refused. The hub can assume exactly one thing: the new connection did not get in.

One caveat on the per-market wording. The ack carries a single pair of numbers, and the WebSocket page does not spell out how that pair relates to a key entitled to several markets. The hub below treats the ack numbers as the budget for the connection it holds, which is the conservative reading; check /docs if you need the exact rule for a multi-market key.

One frame shape across crypto, FX, stocks and commodities#

One hub can serve every service because tick frames are normalized. A crypto tick, an FX tick, a US equity tick and a commodities tick carry the same fields and differ only in class. This is a fixture with the documented field set, chosen to be internally consistent:

{"f":"tick","class":"fx","s":"EURUSD","p":1.16934,"P":85790,"b":1.16925,"B":510300,"a":1.16943,"A":585025,"t":1778019852426}

p and P are the last price and size, b/B the best bid, a/A the best ask, and t is the engine's timestamp for the published value in Unix epoch milliseconds. class is one of cex, dex, fx, us, com; DEX ticks also carry chain. Pool TVL updates arrive as a separate "f":"tvl" frame. The hub routes on class (or tvl) plus s, so a channel name such as md:tick:cex:BTCUSD is the same string the consumer used in its request.

The docs state that multiple subscribes on one connection are fine, and each subscribe frame names one product, so the hub sends one frame per product carrying every symbol for that product. On each subscribe the server first emits one snapshot frame per channel, the last known value from its snapshot cache, in the same shape as a live frame. That single fact is what makes the last-value store trivial: after every (re)subscribe the hub's cache is refilled by the server before the first live tick arrives.

The design#

Redis holds three kinds of keys. md:tick:<product>:<symbol> is a pub/sub channel that carries the raw frame; consumers can PSUBSCRIBE md:tick:cex:* for a whole class. md:last:<product>:<symbol> is a plain string holding the most recent frame with a 24 hour TTL, so a late joiner reads a price before its first live message. md:control is a channel where consumers say what they want, and md:control:reply is where the hub posts refusals.

The hub keeps a wanted map from (product, symbol) to the set of consumers asking for it. A symbol two services both want costs one subscription. When the last consumer for a symbol leaves, or falls silent for 90 seconds, the hub sends an unsubscribe frame so the budget comes back. Requests to the server go through a FIFO of pending requests because error frames do not echo the request they answer; the hub attributes each ack or error to the oldest outstanding request, in order.

hub.py#

# hub.py  Python 3.11+   pip install "websockets>=12" "redis>=5.0.1"
import asyncio
import json
import os
import time
from dataclasses import dataclass, field

import redis.asyncio as redis
import websockets

WS_URL = 'wss://stream.sifting.io/ws/v1'
KEY = os.environ['SIFTING_KEY']
REDIS_URL = os.environ.get('REDIS_URL', 'redis://localhost:6379/0')

CONTROL = 'md:control'          # consumers -> hub: what they want
REPLIES = 'md:control:reply'    # hub -> consumers: refusals
PING_EVERY_S = 30               # docs: 90 s without a client frame closes the socket
CONSUMER_TTL_S = 90             # a consumer silent this long releases its symbols
LAST_TTL_S = 24 * 3600
FATAL = {'auth_failed', 'auth_required', 'auth_timeout', 'max_connections'}


class Fatal(Exception):
    pass


def channel(product, symbol):
    return f'md:tick:{product}:{symbol}'


def last_key(product, symbol):
    return f'md:last:{product}:{symbol}'


@dataclass
class HubState:
    tier: str | None = None
    max_conn: int | None = None
    max_subs: int | None = None
    active_conn: int | None = None
    authed: bool = False
    wanted: dict = field(default_factory=dict)   # (product, symbol) -> {client, ...}
    seen: dict = field(default_factory=dict)     # client -> last heartbeat, monotonic s
    upstream: set = field(default_factory=set)   # (product, symbol) acked upstream
    pending: list = field(default_factory=list)  # requests awaiting ack or error, FIFO
    fatal: str | None = None


def resubscribe_frames(state):
    """One subscribe per product covering the whole union. Runs after every auth ack."""
    state.upstream.clear()
    state.pending.clear()
    by_product = {}
    for product, symbol in state.wanted:
        by_product.setdefault(product, []).append(symbol)
    frames = [{'op': 'subscribe', 'product': p, 'symbols': s} for p, s in by_product.items()]
    state.pending.extend(dict(fr) for fr in frames)
    return frames


def handle_frame(state, frame):
    """Pure: update state, return side effects for the loop to run."""
    f = frame.get('f')

    if f == 'ack':
        op = frame.get('op')
        if op == 'auth':
            state.tier = frame.get('tier')
            state.max_conn = frame.get('max_conn')
            state.max_subs = frame.get('max_subs')
            state.active_conn = frame.get('active_conn')
            state.authed = True
            log = {'kind': 'log', 'msg': f'auth ok tier={state.tier} max_conn={state.max_conn} '
                                        f'max_subs={state.max_subs} active_conn={state.active_conn}'}
            return [log] + [{'kind': 'send', 'frame': fr} for fr in resubscribe_frames(state)]
        if op in ('subscribe', 'unsubscribe'):
            keys = {(frame.get('product'), s) for s in frame.get('symbols', [])}
            state.upstream = (state.upstream | keys) if op == 'subscribe' else (state.upstream - keys)
            if state.pending and state.pending[0]['op'] == op:
                state.pending.pop(0)
        return []

    if f in ('tick', 'tvl'):
        product = 'tvl' if f == 'tvl' else frame.get('class')
        raw = json.dumps(frame, separators=(',', ':'))
        return [{'kind': 'publish', 'channel': channel(product, frame.get('s')), 'data': raw},
                {'kind': 'set', 'key': last_key(product, frame.get('s')), 'data': raw}]

    if f == 'pong':
        return []

    if f == 'error':
        code = frame.get('code')
        if code in FATAL:
            state.fatal = code
            return [{'kind': 'fatal', 'code': code, 'message': frame.get('message'),
                     'limit': frame.get('limit')}]
        if code in ('max_subscriptions', 'unknown_product', 'bad_op') and state.pending:
            req = state.pending.pop(0)          # error frames do not echo the request
            for s in req['symbols']:
                state.wanted.pop((req['product'], s), None)
            return [{'kind': 'reply', 'data': json.dumps({'refused': req, 'code': code,
                     'message': frame.get('message'), 'limit': frame.get('limit')})}]
        return [{'kind': 'log', 'msg': f"error {code}: {frame.get('message')}"}]

    return [{'kind': 'log', 'msg': f'unhandled frame {f}'}]


def release(state, client):
    """Drop a consumer; unsubscribe upstream anything nobody else wants."""
    state.seen.pop(client, None)
    freed = {}
    for key, clients in list(state.wanted.items()):
        clients.discard(client)
        if not clients:
            del state.wanted[key]
            freed.setdefault(key[0], []).append(key[1])
    frames = [{'op': 'unsubscribe', 'product': p, 'symbols': s} for p, s in freed.items()]
    state.pending.extend(dict(fr) for fr in frames)
    return frames


def handle_control(state, msg, now):
    """Pure: one consumer message -> (frames for upstream, replies for consumers)."""
    client, op = msg.get('client'), msg.get('op')
    if not client:
        return [], []
    state.seen[client] = now
    if op == 'leave':
        return release(state, client), []
    if op != 'subscribe':
        return [], []
    product, symbols = msg.get('product'), list(msg.get('symbols', []))
    new = [s for s in symbols if (product, s) not in state.wanted]
    for s in symbols:
        state.wanted.setdefault((product, s), set()).add(client)
    if not new:
        return [], []                       # already streaming, nothing to spend
    if state.max_subs is not None and len(state.wanted) > state.max_subs:
        for s in new:                       # roll back and refuse locally
            del state.wanted[(product, s)]
        reply = {'refused': {'product': product, 'symbols': new}, 'code': 'hub_budget',
                 'message': f'union would exceed max_subs={state.max_subs}',
                 'limit': state.max_subs}
        return [], [reply]
    frame = {'op': 'subscribe', 'product': product, 'symbols': new}
    state.pending.append(dict(frame))
    return [frame], []


def expire_consumers(state, now, ttl=CONSUMER_TTL_S):
    frames = []
    for client, last in list(state.seen.items()):
        if now - last > ttl:
            frames += release(state, client)
    return frames


async def apply(actions, ws, r):
    for a in actions:
        kind = a['kind']
        if kind == 'publish':
            await r.publish(a['channel'], a['data'])
        elif kind == 'set':
            await r.set(a['key'], a['data'], ex=LAST_TTL_S)
        elif kind == 'reply':
            await r.publish(REPLIES, a['data'])
        elif kind == 'send':
            await ws.send(json.dumps(a['frame']))
        else:
            print(a)


async def reader(ws, state, r):
    async for raw in ws:
        await apply(handle_frame(state, json.loads(raw)), ws, r)
        if state.fatal:
            raise Fatal(state.fatal)


async def pinger(ws):
    while True:
        await asyncio.sleep(PING_EVERY_S)
        await ws.send(json.dumps({'op': 'ping'}))


async def control(ws, state, r):
    ps = r.pubsub()
    await ps.subscribe(CONTROL)
    try:
        while True:
            msg = await ps.get_message(ignore_subscribe_messages=True, timeout=5.0)
            now = time.monotonic()
            frames, replies = [], []
            if msg:
                frames, replies = handle_control(state, json.loads(msg['data']), now)
            frames += expire_consumers(state, now)
            if state.authed:                # before the ack, the resubscribe covers it
                for fr in frames:
                    await ws.send(json.dumps(fr))
            for rep in replies:
                await r.publish(REPLIES, json.dumps(rep))
    finally:
        await ps.aclose()


async def main():
    state = HubState()
    r = redis.from_url(REDIS_URL, decode_responses=True)
    backoff = 1
    while True:
        state.authed = False
        try:
            async with websockets.connect(f'{WS_URL}?key={KEY}') as ws:
                async with asyncio.TaskGroup() as tg:
                    tg.create_task(reader(ws, state, r))
                    tg.create_task(pinger(ws))
                    tg.create_task(control(ws, state, r))
        except* Fatal as eg:
            print('fatal, leaving the restart to the supervisor:', eg.exceptions[0])
            raise SystemExit(2)
        except* (websockets.WebSocketException, OSError, TimeoutError) as eg:
            print('upstream dropped:', repr(eg.exceptions[0]))
        backoff = 1 if state.authed else min(backoff * 2, 30)
        print(f'reconnecting in {backoff} s')
        await asyncio.sleep(backoff)


if __name__ == '__main__':
    asyncio.run(main())

A few decisions in that file deserve a sentence each.

The auth ack triggers the resubscribe. The hub never assumes the first frame is the ack; whenever an op: auth ack arrives, on the first connection or the tenth reconnect, resubscribe_frames rebuilds the union as one subscribe frame per product and sends it. Consumers never learn that the upstream dropped. Their Redis subscriptions stay open, the server's snapshot-on-subscribe refills every md:last key, and live ticks resume.

The budget check happens before the round trip. handle_control compares the size of the union after the request against max_subs from the ack and refuses on md:control:reply with the same limit shape the server would use, so consumer code handles one refusal format. If the server still answers max_subscriptions (for example because another connection on the same key is spending the same budget), the hub attributes it to the oldest pending request, drops those symbols from wanted, and relays the server's message and limit.

The pinger is unconditional. The docs are explicit that 90 seconds without a client frame closes the connection and that inbound ticks do not count. A busy stream with no client pings still dies at 90 seconds. Thirty seconds leaves margin under the documented "at least every 60 s".

max_connections is fatal, on purpose. If the hub itself is refused, the key's connection budget is spent elsewhere, and reconnecting every second changes nothing except the log volume. The process exits and a supervisor restarts it after a delay (more on the delay below).

The consumer: subscribe first, then read the snapshot#

Every service runs the same few lines. It restates what it wants every 30 seconds, because the hub forgets a consumer silent for 90 seconds and releases its symbols; that heartbeat doubles as the request, so a service that starts before the hub is picked up as soon as the hub is listening.

# consumer.py  any service: alerting, charting, a pricing cache
import asyncio
import json
import os

import redis.asyncio as redis

REDIS_URL = os.environ.get('REDIS_URL', 'redis://localhost:6379/0')
CLIENT = os.environ.get('CLIENT_ID', 'alerting')
WANT = {'cex': ['BTCUSD', 'ETHUSD'], 'fx': ['EURUSD']}


async def announce(r):
    # the hub forgets a consumer silent for 90 s, so restate the wants every 30 s
    while True:
        for product, symbols in WANT.items():
            await r.publish('md:control', json.dumps(
                {'op': 'subscribe', 'client': CLIENT, 'product': product, 'symbols': symbols}))
        await asyncio.sleep(30)


async def follow(r, product, symbol):
    ps = r.pubsub()
    await ps.subscribe(f'md:tick:{product}:{symbol}')     # 1. listen first
    last_t = 0
    cached = await r.get(f'md:last:{product}:{symbol}')   # 2. then read the snapshot
    if cached:
        frame = json.loads(cached)
        last_t = frame['t']
        yield frame
    async for msg in ps.listen():                          # 3. live, minus anything older
        if msg['type'] != 'message':
            continue
        frame = json.loads(msg['data'])
        if frame['t'] <= last_t:
            continue
        last_t = frame['t']
        yield frame


async def main():
    r = redis.from_url(REDIS_URL, decode_responses=True)
    announcer = asyncio.create_task(announce(r))
    async for frame in follow(r, 'cex', 'BTCUSD'):
        print(frame['s'], 'last', frame.get('p'), 'bid', frame.get('b'),
              'ask', frame.get('a'), 't', frame['t'])
    announcer.cancel()


if __name__ == '__main__':
    asyncio.run(main())

The order in follow is the whole point. Subscribing to the channel before reading md:last means no tick can fall in the gap between the two calls; anything published during the read arrives on the channel. The price of that ordering is a possible duplicate, since the frame stored in md:last may also arrive on the channel. Dropping any frame whose t is not newer than the last one seen handles that, and since t is the engine's timestamp for the published value, an equal or older t never carries a newer price.

The snapshot a late joiner gets is the last value the server had, and that value may be minutes old on a quiet symbol or hours old outside a market's session. The consumer, not the hub, decides what that age means for its purpose, with clock skew and invalid timestamps treated as their own cases. Bid ask spread: how to read a quote and flag a wide or stale one in code covers that judgement.

Tests for the error frames#

Because handle_frame and handle_control touch no socket and no Redis, the error paths run under pytest with no key and no network:

# test_hub.py    pytest -q test_hub.py
import json
import os

os.environ.setdefault('SIFTING_KEY', 'sft_test')     # hub.py reads it at import
from hub import HubState, handle_control, handle_frame


def auth(state, max_subs=5):
    return handle_frame(state, {'f': 'ack', 'op': 'auth', 'tier': 'free',
                                'max_conn': 1, 'max_subs': max_subs, 'active_conn': 1})


def sub(state, client, product, symbols, now=0):
    return handle_control(state, {'op': 'subscribe', 'client': client,
                                  'product': product, 'symbols': symbols}, now)


def test_auth_ack_sets_the_budget():
    s = HubState()
    auth(s)
    assert (s.tier, s.max_conn, s.max_subs, s.active_conn) == ('free', 1, 5, 1)
    assert s.authed


def test_max_connections_is_fatal_and_not_retried_by_the_loop():
    s = HubState()
    actions = handle_frame(s, {'f': 'error', 'code': 'max_connections',
                               'message': 'Tier connection cap reached', 'limit': 1})
    assert s.fatal == 'max_connections'
    assert actions == [{'kind': 'fatal', 'code': 'max_connections',
                        'message': 'Tier connection cap reached', 'limit': 1}]


def test_max_subscriptions_refuses_the_oldest_pending_request():
    s = HubState()
    auth(s, max_subs=1000)
    frames, _ = sub(s, 'charting', 'cex', ['BTCUSD'])
    assert frames == [{'op': 'subscribe', 'product': 'cex', 'symbols': ['BTCUSD']}]
    actions = handle_frame(s, {'f': 'error', 'code': 'max_subscriptions',
                               'message': 'Tier subscription cap reached', 'limit': 1000})
    assert actions[0]['kind'] == 'reply'
    reply = json.loads(actions[0]['data'])
    assert reply['code'] == 'max_subscriptions' and reply['limit'] == 1000
    assert reply['refused']['symbols'] == ['BTCUSD']
    assert s.pending == [] and ('cex', 'BTCUSD') not in s.wanted


def test_hub_refuses_locally_before_spending_a_round_trip():
    s = HubState()
    auth(s, max_subs=2)
    sub(s, 'a', 'cex', ['BTCUSD', 'ETHUSD'])
    frames, replies = sub(s, 'b', 'fx', ['EURUSD'])
    assert frames == []
    assert replies[0]['code'] == 'hub_budget' and replies[0]['limit'] == 2
    assert len(s.wanted) == 2


def test_shared_symbol_costs_one_subscription():
    s = HubState()
    auth(s)
    frames_a, _ = sub(s, 'a', 'cex', ['BTCUSD'])
    frames_b, _ = sub(s, 'b', 'cex', ['BTCUSD', 'ETHUSD'])
    assert frames_a == [{'op': 'subscribe', 'product': 'cex', 'symbols': ['BTCUSD']}]
    assert frames_b == [{'op': 'subscribe', 'product': 'cex', 'symbols': ['ETHUSD']}]


def test_leave_releases_only_what_nobody_else_wants():
    s = HubState()
    auth(s)
    sub(s, 'a', 'cex', ['BTCUSD'])
    sub(s, 'b', 'cex', ['BTCUSD', 'ETHUSD'])
    frames, _ = handle_control(s, {'op': 'leave', 'client': 'b'}, now=1)
    assert frames == [{'op': 'unsubscribe', 'product': 'cex', 'symbols': ['ETHUSD']}]
    assert s.wanted == {('cex', 'BTCUSD'): {'a'}}


def test_reconnect_resubscribes_the_union_in_one_frame_per_product():
    s = HubState()
    auth(s)
    sub(s, 'a', 'cex', ['BTCUSD'])
    sub(s, 'b', 'fx', ['EURUSD'])
    sub(s, 'c', 'cex', ['ETHUSD'])
    s.authed = False                                   # connection dropped
    actions = auth(s)                                  # new connection, new ack
    sends = [a['frame'] for a in actions if a['kind'] == 'send']
    assert sends == [{'op': 'subscribe', 'product': 'cex', 'symbols': ['BTCUSD', 'ETHUSD']},
                     {'op': 'subscribe', 'product': 'fx', 'symbols': ['EURUSD']}]

The message strings in the fixtures are the ones the docs page lists for each code; the limit values are fixtures. Nothing here was run against the live stream, so treat the tests as the contract the hub is written to, and run them before the first deploy.

Gotchas#

The hub gets max_connections on startup. That means the key's connection budget is in use somewhere else: a consumer that was never migrated to Redis, a second hub instance from a blue/green deploy, or the hub's own previous process, whose socket the server has not closed yet. The docs do not say a new connection displaces an older one, so don't design for that. A crashed predecessor sends no client frames, and the server closes a connection after 90 seconds without one, so a supervisor restart delay of at least 90 seconds (RestartSec=90 in a systemd unit) gives a dead predecessor time to be reaped. Log active_conn from the first ack every time; it names the problem.

Redis pub/sub is fire-and-forget. A subscriber that cannot keep up is disconnected by Redis once its output buffer limit is hit, and a service that was down misses everything published while it was gone. For alerting, charting and a pricing cache that is fine, because the md:last key covers the gap. A service that needs every tick, such as a recorder, should read from a Redis Stream (XADD with MAXLEN) that the hub writes alongside the publish, or hold its own connection if the budget allows.

DEX channel names. DEX subscriptions use the chain:PAIR form, for example eth:WETH-USDC, and DEX ticks carry class: "dex" plus a chain field. Whether the s field on a DEX tick includes the chain prefix is not something to assume from this post; check the frame on the docs page and, if s arrives bare, build the channel name from chain and s so it matches what consumers requested.

The key rides in the query string. Any log line that prints the connect URL prints the key. Log the host, never the URL, and keep the WebSocket library's logger above DEBUG in production.

Unsubscribe returns budget, leaving does not. A consumer that exits without sending leave keeps its symbols subscribed for up to 90 seconds until the heartbeat expiry releases them. Services with a clean shutdown path should publish {"op":"leave","client":"..."} so a redeploy of a symbol-heavy service frees its budget immediately rather than a minute and a half later.

Read the docs

Keep reading

Related posts