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 wantInstallConnect to
Posts, likes, follows as JSON (most projects)pip install websocketswss://jetstream2.us-east.bsky.network/subscribe
Raw, signed CBOR commitspip install atprotoThe relay, via FirehoseSubscribeReposClient
Just keyword alerts, no codeNothingKeyword 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 no commit field. Skip them unless you need them.
  • commit.operation – create, update or delete. Deletes have no record, 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 optionally reply, embed and facets.

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.path is collection/rkey, and commit.repo is the DID.
  • Cursors are sequence numbers. Save commit.seq and pass it back with FirehoseSubscribeReposClient(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 use AsyncFirehoseSubscribeReposClient in 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

Download Skyscraper →

Related Reading