Tutorial
Python WebSocket client: stream live market prices with asyncio
Most Python WebSocket tutorials talk to an echo server. This one streams real quotes and trades, survives a dropped connection and leaves you a CSV file that pandas can read.
On this page
Key takeaways
- The `websockets` library (version 14 or newer) and asyncio are all you need; one connection streams BTCUSD and EURUSD together.
- Wait for the `ready` frame before subscribing, and keep your subscriptions in one list so every reconnect can resend them.
- Crypto and forex prices arrive as strings: write them to disk as sent and convert with `float()` only when you calculate.
- Reconnect with exponential backoff plus jitter, and skip snapshot frames that are not newer than the last tick you already stored.
- Protocol ping and pong frames are handled by the library in both directions; you only decide how long silence may last.
A WebSocket client in Python needs four things to be useful for market data: an asyncio connection, a subscribe message sent only after the server's ready frame, a parser that treats prices as strings, and a loop that reconnects when the network drops. The client below does all four in about 130 lines with the websockets library, streams Bitcoin and euro quotes together, and appends every tick to ticks.csv.
It talks to the TickerLayer stream described in the stock WebSocket API guide, which covers the full protocol. The full client was tested against the live stream on 28 September 2026, and every frame in the outputs below was captured from it that day.
- TickerLayer streamwss, JSON frames
- websocketsasyncio connection
- handle()route by type
- TickWriterdedupe, append
- ticks.csvone row per tick
Before you start
- Python 3.9 or newer, and
pip install "websockets>=14"(the code uses the asyncio client andInvalidStatus, both version 14 or later). - A TickerLayer key with WebSocket access. The free tier is REST only; streaming comes with paid plans, and a free account can request a WebSocket trial from the dashboard.
- The key in an environment variable:
export TICKERLAYER_API_KEY=.... Never paste it into the script. - The crypto and forex feeds on your plan, since the client subscribes to
crypto.quotes,crypto.tradesandforex.quotes.
WebSocket authentication and the first frames
Authentication happens once, in the connection URL: the key goes in the apiKey query parameter, URL-encoded. There is no login message and no header to set. The smallest useful client connects, prints the ready frame, subscribes to one pair and prints three frames:
import asyncio
import json
import os
from urllib.parse import quote
import websockets # pip install "websockets>=14"
URL = "wss://stream.tickerlayer.com/?apiKey=" + quote(os.environ["TICKERLAYER_API_KEY"], safe="")
async def main():
async with websockets.connect(URL, compression=None) as ws:
print(await ws.recv()) # the ready frame: nothing may be sent before it
await ws.send(json.dumps({"action": "subscribe", "channels": ["forex.quotes"], "symbols": ["EURUSD"]}))
for _ in range(3):
print(await ws.recv())
asyncio.run(main()){"type":"system","event":"ready","ts":1790591203518}
{"type":"system","event":"subscribed","channels":["forex.quotes"],"symbols":["EURUSD"]}
{"type":"quote","channel":"forex.quotes","asset":"forex","symbol":"EURUSD","bid":"1.137107","ask":"1.137157","bid_size":"44","ask_size":"58","ts":1790591203002}
{"type":"quote","channel":"forex.quotes","asset":"forex","symbol":"EURUSD","bid":"1.137109","ask":"1.137159","bid_size":"44","ask_size":"58","ts":1790591203003}The EURUSD quote frame
{
"type": "quote",1
"channel": "forex.quotes",
"symbol": "EURUSD",
"bid": "1.137107",2
"ask": "1.137157",
"bid_size": "44",
"ask_size": "58",
"ts": 17905912030023
}
typequote,trade,systemorerror. The handler routes on this field first.bidA string.float("1.137107")is fine for arithmetic; keep the string for storage.tsEvent time in Unix milliseconds, UTC. Divide by 1,000 fordatetime.fromtimestamp.
Two details already matter. compression=None turns off per-message compression: the stream sends uncompressed frames anyway, and switching it off in the client guarantees your process never spends time inflating frames on a busy socket. And quote(..., safe="") encodes every reserved character in the key, so a key that contains one cannot break the URL.
WebSocket ping pong: who pings whom
Keepalive runs in both directions, and the websockets library does all of it for you. Understanding it still matters, because the failure mode is a disconnect that looks random.
| Mechanism | Who sends it | Who answers | Your job |
|---|---|---|---|
| Protocol ping from the server | TickerLayer | The library, automatically | Keep the event loop free so the answer goes out in time. |
| Protocol ping from the client | The library, every ping_interval seconds (default 20) | The server | Choose how long to wait: ping_timeout (default 20). |
| JSON ping | You: {"action":"ping"} | The server, with a pong system frame | Optional. Proves the server application is reading you. |
The one you can break is the first. The server expects pongs to its pings; if your code blocks the event loop (a synchronous database call, a slow time.sleep() inside the handler), the pong never leaves, and the server closes the connection with code 4008 after sending HEARTBEAT_TIMEOUT. In an asyncio client, anything slow belongs in await-able code or in a thread via asyncio.to_thread.
A complete WebSocket Python client with reconnects and CSV output
The complete client adds three things to the minimal one: a TickWriter that appends rows and skips replayed values, a handle() function that routes frames by type, and a main() loop that reconnects forever unless the key itself is refused.
import asyncio
import csv
import json
import os
import random
import time
from datetime import datetime, timezone
from pathlib import Path
from urllib.parse import quote
import websockets # pip install "websockets>=14"
API_KEY = os.environ["TICKERLAYER_API_KEY"]
URL = "wss://stream.tickerlayer.com/?apiKey=" + quote(API_KEY, safe="")
CSV_PATH = Path("ticks.csv")
FIELDS = ["time_utc", "ts", "channel", "symbol", "bid", "ask", "price", "size", "snapshot"]
# Everything we want, re-sent after every reconnect.
SUBSCRIPTIONS = [
{"action": "subscribe", "channels": ["crypto.quotes", "crypto.trades"], "symbols": ["BTCUSD"]},
{"action": "subscribe", "channels": ["forex.quotes"], "symbols": ["EURUSD"]},
]
class TickWriter:
"""Append quote and trade frames to a CSV file, skipping replays we already hold."""
def __init__(self, path):
new_file = not path.exists()
self.file = path.open("a", newline="")
self.csv = csv.DictWriter(self.file, fieldnames=FIELDS)
if new_file:
self.csv.writeheader()
self.last_ts = {} # (channel, symbol) -> newest ts written
self.last_flush = time.monotonic()
def write(self, msg):
key = (msg["channel"], msg["symbol"])
ts = int(msg["ts"])
if msg.get("snapshot") and ts <= self.last_ts.get(key, 0):
return False # a snapshot of something written before the reconnect
self.last_ts[key] = max(ts, self.last_ts.get(key, 0))
self.csv.writerow({
"time_utc": datetime.fromtimestamp(ts / 1000, tz=timezone.utc).isoformat(timespec="milliseconds"),
"ts": ts,
"channel": msg["channel"],
"symbol": msg["symbol"],
# Numbers arrive as strings on crypto, forex and stocks: store them as sent.
"bid": msg.get("bid", ""),
"ask": msg.get("ask", ""),
"price": msg.get("price", ""),
"size": msg.get("size", ""),
"snapshot": 1 if msg.get("snapshot") else 0,
})
if time.monotonic() - self.last_flush > 1:
self.file.flush()
self.last_flush = time.monotonic()
return True
def close(self):
self.file.close()
def handle(msg, writer):
kind = msg.get("type")
if kind in ("quote", "trade"):
if not writer.write(msg):
return
tag = " (snapshot)" if msg.get("snapshot") else ""
if kind == "quote":
bid, ask = float(msg["bid"]), float(msg["ask"])
bps = (ask - bid) / ((ask + bid) / 2) * 10_000
print(f"{msg['symbol']:<7} quote {msg['bid']} / {msg['ask']} ({bps:.2f} bps){tag}")
else:
print(f"{msg['symbol']:<7} trade {msg['price']} x {msg['size']}{tag}")
elif kind == "error":
# A failed subscribe does not close the socket; fix the request, do not reconnect.
print("error:", msg.get("code"), msg.get("message"), msg.get("rejected", ""))
elif kind == "system" and msg.get("event") == "disconnect":
print("server is closing the stream:", msg.get("code"))
elif kind == "system":
print("system:", msg.get("event"), msg.get("channels", ""), msg.get("symbols", ""))
async def stream(writer):
# compression=None: the stream is uncompressed, so skip the negotiation entirely.
# ping_interval/ping_timeout: our own keepalive; the library also answers the
# server's protocol pings on its own.
async with websockets.connect(
URL, compression=None, open_timeout=10, ping_interval=20, ping_timeout=20
) as ws:
first = json.loads(await asyncio.wait_for(ws.recv(), timeout=10))
if first.get("event") != "ready":
raise RuntimeError(f"expected the ready frame, got {first}")
for sub in SUBSCRIPTIONS:
await ws.send(json.dumps(sub))
async for raw in ws:
handle(json.loads(raw), writer)
print(f"closed cleanly with code {ws.close_code}")
async def main():
writer = TickWriter(CSV_PATH)
backoff = 1.0
try:
while True:
started = time.monotonic()
try:
await stream(writer)
except websockets.ConnectionClosed as exc:
code = exc.rcvd.code if exc.rcvd else 1006 # 1006: no close frame at all
print(f"connection lost, close code {code}")
except websockets.InvalidStatus as exc:
status = exc.response.status_code
if status in (401, 403):
raise SystemExit(f"HTTP {status} on upgrade: check the key and the plan")
print(f"HTTP {status} on upgrade, will retry")
except (OSError, asyncio.TimeoutError) as exc:
print(f"network error: {exc!r}")
if time.monotonic() - started > 30:
backoff = 1.0 # that connection was healthy: restart the ladder
delay = backoff + random.uniform(0, backoff / 2)
print(f"reconnecting in {delay:.1f}s")
await asyncio.sleep(delay)
backoff = min(backoff * 2, 30.0)
finally:
writer.close()
if __name__ == "__main__":
try:
asyncio.run(main())
except KeyboardInterrupt:
passRun it with python stream_to_csv.py and stop it with Ctrl+C. For the output below we cut the network connection five seconds in, to show the recovery path:
system: subscribed ['crypto.quotes', 'crypto.trades'] ['BTCUSD']
system: subscribed ['forex.quotes'] ['EURUSD']
BTCUSD trade 82991.98 x 3e-8
BTCUSD trade 83028.53 x 0.00057
BTCUSD quote 83028.53 / 83028.54 (0.00 bps)
EURUSD quote 1.137258 / 1.137297 (0.34 bps)
BTCUSD trade 83028.54 x 0.00012
EURUSD quote 1.137256 / 1.137296 (0.35 bps)
...
BTCUSD trade 83027.6 x 0.006021
connection lost, close code 1006
reconnecting in 1.2s
system: subscribed ['crypto.quotes', 'crypto.trades'] ['BTCUSD']
BTCUSD trade 82980.63 x 3e-8
system: subscribed ['forex.quotes'] ['EURUSD']
BTCUSD trade 83024.99 x 0.00042
BTCUSD quote 83024.99 / 83025 (0.00 bps)
EURUSD quote 1.137285 / 1.137295 (0.09 bps)The drop surfaced as close code 1006, which means the connection died without a close frame: exactly what a Wi-Fi switch or a NAT timeout looks like. Close 1012 (a planned release, announced by a SERVER_RESTART system frame) takes the same path. Note the BTC trades: 3e-8 is a legitimate size, which is one more reason to keep the strings as sent. The WebSocket reconnect guide walks through every close code and how fast to come back from each.
How the reconnect loop behaves
Each failure doubles the base delay up to 30 seconds and adds up to half of it again at random. The jitter matters when many clients drop at once: without it they all return in the same second. The ladder resets only after a connection has stayed up for 30 seconds, so a server that accepts and immediately drops you cannot pull the client into a tight loop.
| Failure in a row | Base delay | Actual wait |
|---|---|---|
| 1st | 1 s | 1 to 1.5 s |
| 2nd | 2 s | 2 to 3 s |
| 3rd | 4 s | 4 to 6 s |
| 4th | 8 s | 8 to 12 s |
| 5th | 16 s | 16 to 24 s |
| 6th and later | 30 s | 30 to 45 s |
Two statuses stop the loop instead: HTTP 401 and 403 on the upgrade. A bad key or a plan without streaming will never succeed on retry, and hammering the endpoint with it only fills your logs.
Why snapshot frames need deduplication
After every subscribe, the stream can replay the last recent quote or trade per symbol, marked "snapshot": true, so a screen has a value immediately. After a reconnect that replay may be a tick you already wrote before the drop. TickWriter keeps the newest ts per channel and symbol and skips a snapshot unless it is newer. Live frames are never skipped, because two genuine trades can share a millisecond.
In this run every row has snapshot 0: crypto and forex trade continuously, so they rarely replay. Add US:KO on stocks.quotes during the US session and its first row will usually be a snapshot. The snapshot rules explain when a replay exists and how "snapshot": false turns it off.
time_utc,ts,channel,symbol,bid,ask,price,size,snapshot
2026-09-28T10:53:33.681+00:00,1790592813681,crypto.trades,BTCUSD,,,82991.98,3e-8,0
2026-09-28T10:53:33.640+00:00,1790592813640,crypto.trades,BTCUSD,,,83028.53,0.00057,0
2026-09-28T10:53:33.791+00:00,1790592813791,crypto.quotes,BTCUSD,83028.53,83028.54,,,0
2026-09-28T10:53:33.002+00:00,1790592813002,forex.quotes,EURUSD,1.137258,1.137297,,,0Look at the first two rows: the second trade has an earlier ts than the first. Rows are in arrival order, and trades from an aggregated feed do not always arrive in event order, so sort by ts before you compute anything time-based.
Read the CSV back with pandas
The file is ready for analysis as soon as it has a few rows. This script drops snapshot rows, computes each quote's spread in basis points and resamples BTC trades into five-second bars:
import pandas as pd
ticks = pd.read_csv("ticks.csv")
ticks["time_utc"] = pd.to_datetime(ticks["time_utc"])
live = ticks[ticks["snapshot"] == 0]
quotes = live[live["channel"].str.endswith(".quotes")].copy()
quotes["mid"] = (quotes["bid"] + quotes["ask"]) / 2
quotes["spread_bps"] = (quotes["ask"] - quotes["bid"]) / quotes["mid"] * 10_000
print(quotes.groupby("symbol")["spread_bps"].describe()[["count", "mean", "min", "max"]].round(3))
trades = live[live["channel"] == "crypto.trades"].set_index("time_utc").sort_index()
bars = trades["price"].resample("5s").ohlc().dropna()
print(bars.tail(3)) count mean min max
symbol
BTCUSD 32.0 0.001 0.001 0.001
EURUSD 38.0 0.297 0.088 0.440
open high low close
time_utc
2026-09-28 10:53:30+00:00 83028.53 83042.63 82991.98 82994.27
2026-09-28 10:53:35+00:00 82994.27 83045.68 82980.63 83040.76
2026-09-28 10:53:40+00:00 82980.63 83040.76 82980.00 82980.01From here the file feeds a backtest, a spread monitor or the signal logic in the Python trading bot tutorial. For history older than your recording, load bars over REST instead; WebSocket vs REST covers which job belongs to which. The live pages for BTCUSD and EURUSD show the same instruments with a public delay.
Production notes
- Never log the URLThe key sits in the query string. Log the host, never the full URL, and keep the key in the environment.
- Keep handle() fastAnything slow blocks pongs and ends in close 4008. Push frames onto an
asyncio.Queueand write from a separate task when volume grows. - One socket per processPlans cap connections (one per feed on Individual). Add symbols to the existing socket instead of opening another.
- Mind the gapA reconnect restores the latest value, not the ticks you missed. Backfill bars over REST if the gap matters.
The websockets library can also reconnect on its own: async for ws in websockets.connect(URL) retries with its own backoff. The explicit loop above is longer but makes every decision visible, including which statuses should never be retried.
Questions
How do I create a WebSocket client in Python?
Install websockets, open the connection with async with websockets.connect(url), then await ws.send() and await ws.recv() or iterate with async for message in ws. Run it with asyncio.run().
Which Python WebSocket library should I use?
websockets is the standard asyncio choice and handles ping, pong and close codes for you. websocket-client is a synchronous, callback-based alternative that suits scripts without asyncio.
How does WebSocket ping pong work in Python?
The websockets library answers server pings automatically and sends its own every ping_interval seconds, closing the connection if no pong arrives within ping_timeout. Blocking the event loop is what breaks it.
How do I pass an API key to a WebSocket in Python?
For TickerLayer, URL-encode it into the apiKey query parameter: "wss://stream.tickerlayer.com/?apiKey=" + quote(key, safe=""). Read the key from an environment variable and never log the full URL.
Why does my Python WebSocket disconnect after a while?
Usually a network reset (close 1006), a server release (1012) or missed heartbeats because the event loop was blocked (4008). Reconnect with backoff, resubscribe, and move slow work out of the message handler.