jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168# /// script# requires-python = ">=3.12"# dependencies = ["zstandard", "requests"]# ///"""backfill one collection from a stream instance's archive into sqlite.
the consumer-side recipe for network.bsky.jetstream.planSnapshot:plan (collection filter) -> fetch planned blocks/segments -> decompress(plain zstd frames, no dictionary) -> columnar decode -> keep matchingrows -> insert. prints a timing/throughput report at the end.
block fetches run concurrently (--workers, default 12): sparse collectionstouch one block per ~4096 events, so a serial fetch is pure round-triplatency. decode/insert stays on the main thread; sqlite needs no locking.
usage: uv run backfill-collection-sqlite.py \ --host http://127.0.0.1:8080 \ --collection place.stream.chat.message \ --db streamplace-chat.db"""
from __future__ import annotations
import argparseimport jsonimport sqlite3import structimport timefrom concurrent.futures import ThreadPoolExecutor
import requestsimport zstandard
def decode_block(raw: bytes, wanted: str) -> list[tuple]: n = struct.unpack_from("<I", raw, 0)[0] if n == 0: return [] o = 4 seqs = struct.unpack_from(f"<{n}Q", raw, o); o += 8 * n wits = struct.unpack_from(f"<{n}q", raw, o); o += 8 * n o += 8 * n # indexed_at o += n # kind col_lens = struct.unpack_from(f"<{n}B", raw, o); o += n did_lens = struct.unpack_from(f"<{n}H", raw, o); o += 2 * n rkey_lens = struct.unpack_from(f"<{n}B", raw, o); o += n rev_lens = struct.unpack_from(f"<{n}B", raw, o); o += n pay_lens = struct.unpack_from(f"<{n}I", raw, o); o += 4 * n
def blob(lens): nonlocal o vals, at = [], o for ln in lens: vals.append(raw[at : at + ln]); at += ln o = at return vals
cols = blob(col_lens) dids = blob(did_lens) rkeys = blob(rkey_lens) blob(rev_lens) pays = blob(pay_lens) want = wanted.encode() return [ (seqs[i], wits[i], dids[i].decode(), rkeys[i].decode(), pays[i]) for i in range(n) if cols[i] == want ]
def main() -> None: ap = argparse.ArgumentParser() ap.add_argument("--host", default="http://127.0.0.1:8080") ap.add_argument("--collection", required=True) ap.add_argument("--db", required=True) ap.add_argument("--workers", type=int, default=12) args = ap.parse_args()
db = sqlite3.connect(args.db) db.execute("pragma journal_mode=wal") db.execute( "create table if not exists records(" "seq integer primary key, witnessed_at integer, did text, rkey text, record blob)" )
s = requests.Session() dctx = zstandard.ZstdDecompressor() t0 = time.monotonic() plan = s.post( f"{args.host}/xrpc/network.bsky.jetstream.planSnapshot", json={"collections": [args.collection]}, timeout=60, ).json() segs = plan["segments"] t_plan = time.monotonic() - t0
def fetch_block(job): name, bi = job resp = s.get( f"{args.host}/xrpc/network.bsky.jetstream.getBlock", params={"segment": name, "blockIndex": bi}, timeout=120, ) resp.raise_for_status() return resp.content
pool = ThreadPoolExecutor(max_workers=args.workers) rows = fetched_bytes = blocks_read = 0 for i, seg in enumerate(segs): if seg["mode"] == "blocks": jobs = [ (seg["name"], bi) for r in seg["blocks"] for bi in range(r["first"], r["last"] + 1) ] for content in pool.map(fetch_block, jobs): fetched_bytes += len(content) blocks_read += 1 batch = decode_block(dctx.decompress(content), args.collection) db.executemany("insert or replace into records values(?,?,?,?,?)", batch) rows += len(batch) else: # whole segment resp = s.get( f"{args.host}/xrpc/network.bsky.jetstream.getSegment", params={"name": seg["name"]}, timeout=600, ) resp.raise_for_status() body = resp.content fetched_bytes += len(body) block_index_offset = struct.unpack_from("<Q", body, 90)[0] block_count = struct.unpack_from("<I", body, 14)[0] for bi in range(block_count): off, csize = struct.unpack_from("<QI", body, block_index_offset + bi * 52) frame = body[off + 8 : off + 8 + csize] blocks_read += 1 batch = decode_block(dctx.decompress(frame), args.collection) db.executemany("insert or replace into records values(?,?,?,?,?)", batch) rows += len(batch) if (i + 1) % 100 == 0: db.commit() el = time.monotonic() - t0 print( f"{i+1}/{len(segs)} segments · {rows:,} rows · " f"{fetched_bytes/1e9:.2f} GB fetched · {el:.0f}s · {fetched_bytes/1e6/el:.0f} MB/s", flush=True, ) db.commit()
el = time.monotonic() - t0 print(json.dumps({ "collection": args.collection, "planned_through_seq": plan["plannedThroughSeq"], "plan_seconds": round(t_plan, 2), "segments": len(segs), "blocks_read": blocks_read, "rows": rows, "fetched_gb": round(fetched_bytes / 1e9, 2), "elapsed_seconds": round(el, 1), "mb_per_second": round(fetched_bytes / 1e6 / el, 1), }, indent=2))
if __name__ == "__main__": main()