To read the Bluesky firehose in Python, connect to Jetstream with the websockets library and parse each message as JSON. It takes about twenty lines, needs no account or API key, and you can filter to posts only with ?wantedCollections=app.bsky.feed.post.
This is the hands-on version. If you want the background first – what the firehose is and when Jetstream is the wrong choice – read our firehose streaming guide and firehose vs Jetstream comparison. Everything below assumes Python 3.10 or later.
Quick Answer
| You want | Install | Connect to |
|---|---|---|
| Posts, likes, follows as JSON (most projects) | pip install websockets | wss://jetstream2.us-east.bsky.network/subscribe |
| Raw, signed CBOR commits | pip install atproto | The relay, via FirehoseSubscribeReposClient |
| Just keyword alerts, no code | Nothing | Keyword Alerts |
Step 1: Connect to Jetstream
Bluesky runs four public Jetstream instances: jetstream1.us-east.bsky.network, jetstream2.us-east.bsky.network, jetstream1.us-west.bsky.network and jetstream2.us-west.bsky.network. They carry the same events. Pick the one closest to your server and keep another in reserve.
Install the one dependency and run this:
# pip install websockets
import asyncio
import json
import websockets
URL = (
"wss://jetstream2.us-east.bsky.network/subscribe"
"?wantedCollections=app.bsky.feed.post"
)
async def main() -> None:
async with websockets.connect(URL) as ws:
async for message in ws:
event = json.loads(message)
if event.get("kind") != "commit":
continue
commit = event["commit"]
if commit["operation"] != "create":
continue
text = commit["record"].get("text", "")
print(event["did"], text[:80].replace("\n", " "))
asyncio.run(main())
You will see posts scroll past immediately. The wantedCollections parameter is server-side filtering: Jetstream drops likes, follows and everything else before it reaches you, which is most of the bandwidth. Repeat the parameter to ask for more than one collection.
Step 2: Understand the Event
Every message is a JSON object. A new post looks like this:
{
"did": "did:plc:abc123...",
"time_us": 1756480000123456,
"kind": "commit",
"commit": {
"rev": "3lx...",
"operation": "create",
"collection": "app.bsky.feed.post",
"rkey": "3lxabc2def3gh",
"record": {
"$type": "app.bsky.feed.post",
"text": "Trying out Jetstream from Python",
"langs": ["en"],
"createdAt": "2026-08-29T14:00:00.000Z"
},
"cid": "bafyrei..."
}
}
The fields you will actually use:
did– the author's permanent identifier. Handles can change; DIDs do not.time_us– a Unix timestamp in microseconds. This is also your cursor (see Step 4).kind–"commit"for record changes. You will also see"identity"and"account"events, which have nocommitfield. Skip them unless you need them.commit.operation–create,updateordelete. Deletes have norecord, so never assume it is there.commit.rkey– the record key, which you need for the post URL.commit.record– the post itself:text,createdAt,langs, and optionallyreply,embedandfacets.
Replies are posts too. If you only want top-level posts, skip records that have a reply key.
Step 3: Filter by Keyword and Build a Post URL
Jetstream cannot filter by text, so keyword matching happens in your code. Use a word-boundary regex rather than in: a plain substring check for "art" matches "start", "party" and "smart".
import re
PATTERN = re.compile(r"\b(skyscraper|jetstream|atproto)\b", re.IGNORECASE)
def post_url(did: str, rkey: str) -> str:
return f"https://bsky.app/profile/{did}/post/{rkey}"
def handle(event: dict) -> None:
if event.get("kind") != "commit":
return
commit = event["commit"]
if commit.get("operation") != "create":
return
text = commit.get("record", {}).get("text", "")
match = PATTERN.search(text)
if match:
print(f"[{match.group(0).lower()}] {post_url(event['did'], commit['rkey'])}")
print(" " + text.replace("\n", " ")[:140])
The post URL is https://bsky.app/profile/{did}/post/{rkey}. bsky.app accepts a DID where a handle normally goes, so you do not need to resolve anything. The equivalent AT URI, if you are passing it to the API, is at://{did}/app.bsky.feed.post/{rkey}.
Two practical notes. Filter on record["langs"] if you only want one language; it is set by the posting client and is usually present. And keep the handler fast: do slow work (database writes, HTTP calls) on a queue, not inline, or you will fall behind during busy periods.
Step 4: Reconnect Without Losing Events
Connections drop. Jetstream restarts, your network blips, your laptop sleeps. Without a cursor, everything posted while you were disconnected is simply gone from your point of view.
The fix: save the time_us of the last event you handled, and reconnect with &cursor=<time_us>. Rewind it a few seconds first, so an event that arrived but was not fully processed is replayed rather than skipped.
import asyncio
import json
import random
from pathlib import Path
import websockets
BASE = "wss://jetstream2.us-east.bsky.network/subscribe"
CURSOR_FILE = Path("cursor.txt")
REWIND_US = 5 * 1_000_000 # five seconds, in microseconds
def load_cursor() -> int | None:
try:
return int(CURSOR_FILE.read_text()) - REWIND_US
except (FileNotFoundError, ValueError):
return None
def build_url(cursor: int | None) -> str:
url = f"{BASE}?wantedCollections=app.bsky.feed.post"
if cursor:
url += f"&cursor={cursor}"
return url
async def run() -> None:
delay = 1
while True:
try:
async with websockets.connect(build_url(load_cursor())) as ws:
print("connected")
delay = 1
last_saved = 0.0
async for message in ws:
event = json.loads(message)
handle(event) # from the previous snippet
now = asyncio.get_running_loop().time()
if now - last_saved > 5:
CURSOR_FILE.write_text(str(event["time_us"]))
last_saved = now
except (websockets.exceptions.WebSocketException, OSError, TimeoutError) as exc:
print(f"disconnected: {exc!r}")
wait = delay + random.random()
print(f"reconnecting in {wait:.1f}s")
await asyncio.sleep(wait)
delay = min(delay * 2, 60)
asyncio.run(run())
What this does, and why:
- Saves the cursor every five seconds, not on every event. Writing a file per event is slow at firehose rates.
- Rewinds five seconds on reconnect. You will see a few events twice, so make your processing idempotent: key stored posts on
did+rkey, and an insert-or-ignore handles the rest. - Backs off exponentially with jitter, capped at 60 seconds. When a Jetstream instance restarts, every client reconnects at once; jitter keeps you out of the stampede.
- Relies on websockets' keepalive pings (on by default) to notice a dead connection that never sent a close frame.
Because the cursor is a timestamp rather than a server-specific sequence number, you can generally reconnect to a different Jetstream instance with the same cursor. Replay only covers a bounded window, though. If your consumer is down for a long time, plan for a gap.
Step 5: The Raw Firehose with the atproto SDK
If you need the signed, full-fidelity stream (com.atproto.sync.subscribeRepos), use the atproto package. It handles the WebSocket framing and CBOR decoding; you unpack each commit's CAR blocks to get the records.
# pip install atproto
from atproto import (
CAR,
FirehoseSubscribeReposClient,
models,
parse_subscribe_repos_message,
)
client = FirehoseSubscribeReposClient()
def on_message(message) -> None:
commit = parse_subscribe_repos_message(message)
if not isinstance(commit, models.ComAtprotoSyncSubscribeRepos.Commit):
return # identity, account and other event types
if not commit.blocks:
return
car = CAR.from_bytes(commit.blocks)
for op in commit.ops:
if op.action != "create" or not op.cid:
continue
if not op.path.startswith("app.bsky.feed.post/"):
continue
record = car.blocks.get(op.cid)
if record and "text" in record:
rkey = op.path.split("/", 1)[1]
print(f"https://bsky.app/profile/{commit.repo}/post/{rkey}")
print(" " + record["text"][:120].replace("\n", " "))
client.start(on_message)
Differences from Jetstream you will notice straight away:
- No server-side filtering. You receive every collection from every repo and discard most of it. That is why the path check comes first.
- The post URL comes from the op path.
op.pathiscollection/rkey, andcommit.repois the DID. - Cursors are sequence numbers. Save
commit.seqand pass it back withFirehoseSubscribeReposClient(models.ComAtprotoSyncSubscribeRepos.Params(cursor=saved_seq)). - It is CPU-heavy.
client.start()blocks and runs your handler for every commit on the network. If you do real work per event, hand it off to worker processes, or useAsyncFirehoseSubscribeReposClientin an asyncio program.
If none of that is something you need, go back to Jetstream. It is the right default.
Rate Limits and Writing Back
Reading the stream is not rate-limited the way the XRPC API is: one connection, everything pushed to you. The limits apply the moment your script acts on what it reads – replying, liking, following, or calling getProfile for every author it sees. A keyword bot that replies to every match on a busy term will hit write limits within the hour. Our rate limits guide has the numbers and the backoff pattern.
The usual fix is to batch: collect DIDs from the stream and look up profiles in groups, rather than one request per event.
Frequently Asked Questions
Do I need a Bluesky account or API key to read the firehose?
No. Both Jetstream and the raw firehose are public, unauthenticated WebSocket streams. You only need credentials if your program also posts, likes or replies, and those writes are subject to the normal API rate limits.
Which Python library should I use for Jetstream?
Jetstream is plain JSON over a WebSocket, so the websockets library plus the standard json module is enough. You do not need an AT Protocol SDK until you want the raw CBOR firehose or want to write back to Bluesky.
How do I get the link to a post from a firehose event?
Combine the event's DID and the commit's rkey: https://bsky.app/profile/{did}/post/{rkey}. The DID works in place of a handle, so you never need an extra lookup to build a working link.
What happens to events I miss while my script is down?
Save the time_us of the last event you processed and reconnect with ?cursor= set to that value minus a few seconds. Jetstream replays from there, within its retention window. Expect a few duplicates from the rewind and handle them idempotently.
Is the raw firehose too slow for Python?
It depends on what you do per event. Decoding every CAR block on the full network is real CPU work, and a single Python process doing heavy work per event will fall behind at peak. If you do not need signatures, Jetstream avoids the problem entirely.
Skip the Plumbing
- Keyword Alerts – an email when a keyword appears on Bluesky, with no consumer to run
- Trending hashtags leaderboard – live hashtag counts, already aggregated
- Skyscraper – a free native Bluesky client for iPhone, iPad and Mac, when you would rather read the posts than parse them