在本示例中,您将消费公共 Bluesky Jetstream ↗ firehose——网络上每条帖子、点赞、转发、关注和屏蔽的实时 WebSocket stream——并将其写入 R2 Data Catalog 作为可查询的 Apache Iceberg 表。
您将学习一个核心 Pipelines 模式:将所有事件发送到同一个 stream,然后使用包含多条 SQL 语句的单个 pipeline 将("扇出")该 stream 路由到多个目标表(每种事件类型一个),而无需为每种类型运行单独的 pipeline。
flowchart TD
A[Bluesky Jetstream WebSocket] --> B[Durable Object]
B -->|send| C[bsky_events_stream]
C --> D[bsky_pipeline]
D --> E[bsky_post]
D --> F[bsky_like]
D --> G[bsky_repost]
D --> H[bsky_follow]
D --> I[bsky_block]
- 注册 Cloudflare 账户 ↗。
- 安装
Node.js↗。
Node.js 版本管理器
使用 Volta ↗ 或 nvm ↗ 等 Node 版本管理器,以避免权限问题并切换 Node.js 版本。本指南后续将介绍的 Wrangler 需要 Node 版本 16.17.0 或更高。
您还需要一个具有 Admin Read & Write(管理员读取和写入) 权限的 R2 API 令牌,其中包括 R2 Data Catalog 和 R2 SQL 访问权限。您将为每个 sink 传递此令牌。不需要 Bluesky 账户或 API 密钥,因为 Jetstream 是公开且无需身份验证的。
通过运行以下命令创建新的 Worker 项目:
npm create cloudflare@latest -- bluesky-pipelineyarn create cloudflare bluesky-pipelinepnpm create cloudflare@latest bluesky-pipeline进行设置时,请选择以下选项:
- 对于 What would you like to start with?,选择
Hello World example。 - 对于 Which template would you like to use?,选择
Worker only。 - 对于 Which language do you want to use?,选择
TypeScript。 - 对于 Do you want to use git for version control?,选择
Yes。 - 对于 Do you want to deploy your application?,选择
No(部署前我们还会做一些修改)。
进入新项目目录:
cd bluesky-pipeline后面使用的 pipelines 命令需要 Wrangler v4 或更高版本。如果您的项目是用旧版本创建的,请立即更新:
npm i -D wrangler@latestyarn add -D wrangler@latestpnpm add -D wrangler@latestbun add -d wrangler@latestStream 只有一个 schema。由于每种事件类型都流经同一个 stream,schema 是您在所有类型中想要的字段的并集。后面的每条路由语句仅选择与其目标相关的列。
在项目根目录创建 schema.json 文件:
{
"fields": [
{ "name": "event_id", "type": "string", "required": true },
{ "name": "event_type", "type": "string", "required": true },
{ "name": "did", "type": "string", "required": false },
{ "name": "operation", "type": "string", "required": false },
{ "name": "event_time", "type": "timestamp", "required": false },
{ "name": "created_at", "type": "string", "required": false },
{ "name": "text", "type": "string", "required": false },
{ "name": "langs", "type": "string", "required": false },
{ "name": "subject_uri", "type": "string", "required": false },
{ "name": "subject_did", "type": "string", "required": false }
]
}您的 sink 将数据写入 R2 Data Catalog 中的 Iceberg 表,因此您需要一个启用了 catalog 的存储桶。
创建名为 bluesky-pipeline 的 R2 存储桶:
npx wrangler r2 bucket create bluesky-pipelineyarn wrangler r2 bucket create bluesky-pipelinepnpm wrangler r2 bucket create bluesky-pipeline在存储桶上启用 R2 Data Catalog:
npx wrangler r2 bucket catalog enable bluesky-pipelineyarn wrangler r2 bucket catalog enable bluesky-pipelinepnpm wrangler r2 bucket catalog enable bluesky-pipeline运行此命令时,记下 Warehouse name。您将需要它来使用 R2 SQL 查询数据。
首先,从 schema 文件创建 stream:
npx wrangler pipelines streams create bsky_events_stream --schema-file schema.jsonyarn wrangler pipelines streams create bsky_events_stream --schema-file schema.jsonpnpm wrangler pipelines streams create bsky_events_stream --schema-file schema.json在输出中记下 stream ID。您将在下一步中使用它来配置 Worker 绑定。
接下来,为每个目标表创建一个 sink。每个 sink 写入 R2 Data Catalog 中自己的 Iceberg 表。将 YOUR_CATALOG_TOKEN 替换为您的 R2 API 令牌。
for t in post like repost follow block; do
npx wrangler pipelines sinks create bsky_${t}_sink \
--type r2-data-catalog \
--bucket bluesky-pipeline \
--namespace bluesky \
--table bsky_${t} \
--catalog-token YOUR_CATALOG_TOKEN \
--roll-interval 60
done现在创建一个 SQL 包含多条 INSERT 语句(每条路由一条)的 pipeline。每条语句按 event_type 过滤 stream,并仅投影与其表相关的列。
创建 fanout.sql 文件:
INSERT INTO bsky_post_sink
SELECT event_id, did, operation, event_time, created_at, text, langs
FROM bsky_events_stream WHERE event_type = 'post';
INSERT INTO bsky_like_sink
SELECT event_id, did, operation, event_time, created_at, subject_uri
FROM bsky_events_stream WHERE event_type = 'like';
INSERT INTO bsky_repost_sink
SELECT event_id, did, operation, event_time, created_at, subject_uri
FROM bsky_events_stream WHERE event_type = 'repost';
INSERT INTO bsky_follow_sink
SELECT event_id, did, operation, event_time, created_at, subject_did
FROM bsky_events_stream WHERE event_type = 'follow';
INSERT INTO bsky_block_sink
SELECT event_id, did, operation, event_time, created_at, subject_did
FROM bsky_events_stream WHERE event_type = 'block';从文件创建 pipeline:
npx wrangler pipelines create bsky_pipeline --sql-file fanout.sqlyarn wrangler pipelines create bsky_pipeline --sql-file fanout.sqlpnpm wrangler pipelines create bsky_pipeline --sql-file fanout.sql一个 pipeline 写入五个表。要添加新的事件类型,添加一个 sink 和一条匹配的 INSERT 语句。Pipeline SQL 创建后无法修改,因此您需要删除并重新创建 pipeline 来更改它。有关更多信息,请参阅将一个 stream 路由到多个表。
添加 stream 绑定、用于保持 WebSocket 连接的 Durable Object,以及保持消费者运行的 cron 触发器。将 <STREAM_ID> 替换为步骤 4 中的 stream ID。
{
"$schema": "./node_modules/wrangler/config-schema.json",
"name": "bluesky-pipeline",
"main": "src/index.ts",
// Set this to today's date
"compatibility_date": "2026-08-17",
"pipelines": [
{
"binding": "BSKY_STREAM",
"stream": "<STREAM_ID>"
}
],
"durable_objects": {
"bindings": [
{
"name": "JETSTREAM",
"class_name": "JetstreamConsumer"
}
]
},
"migrations": [
{
"tag": "v1",
"new_sqlite_classes": [
"JetstreamConsumer"
]
}
],
"triggers": {
"crons": [
"*/2 * * * *"
]
}
}name = "bluesky-pipeline"
main = "src/index.ts"
# Set this to today's date
compatibility_date = "2026-08-17"
[[pipelines]]
binding = "BSKY_STREAM"
stream = "<STREAM_ID>"
[[durable_objects.bindings]]
name = "JETSTREAM"
class_name = "JetstreamConsumer"
[[migrations]]
tag = "v1"
new_sqlite_classes = ["JetstreamConsumer"]
[triggers]
crons = ["*/2 * * * *"]Durable Object 是长期 WebSocket 的理想宿主。它在 socket 打开时保持驻留,alarm 在连接断开时重新连接。缓冲传入事件并批量 send() 到 stream,以保持在每次请求 5 MB 的限制以下。持久化 Jetstream time_us 游标,以便重连时无间隙恢复。
将 src/index.ts 的内容替换为以下内容:
import { DurableObject } from "cloudflare:workers";
// Jetstream collection -> our short event_type. Only these are kept.
const COLLECTION_TO_TYPE = {
"app.bsky.feed.post": "post",
"app.bsky.feed.like": "like",
"app.bsky.feed.repost": "repost",
"app.bsky.graph.follow": "follow",
"app.bsky.graph.block": "block",
};
const WANTED = Object.keys(COLLECTION_TO_TYPE);
const JETSTREAM_URL = "https://jetstream2.us-east.bsky.network/subscribe";
const FLUSH_MAX = 500; // rows per send()
const FLUSH_MS = 1000; // flush at least once per second
const RECONNECT_MS = 15000;
// Flatten one Jetstream message into a unified stream row, or null to skip.
function toRow(ev) {
if (ev?.kind !== "commit" || !ev.commit) return null;
const c = ev.commit;
const event_type = COLLECTION_TO_TYPE[c.collection];
if (!event_type) return null;
const r = c.record ?? {};
const subject = r.subject;
return {
event_id: `${ev.did}/${c.collection}/${c.rkey}`,
event_type,
did: ev.did ?? null,
operation: c.operation ?? null,
event_time:
typeof ev.time_us === "number"
? new Date(ev.time_us / 1000).toISOString()
: null,
created_at: typeof r.createdAt === "string" ? r.createdAt : null,
text: event_type === "post" && typeof r.text === "string" ? r.text : null,
langs:
event_type === "post" && Array.isArray(r.langs)
? r.langs.join(",")
: null,
subject_uri: typeof subject === "object" ? (subject?.uri ?? null) : null,
subject_did: typeof subject === "string" ? subject : null,
};
}
export class JetstreamConsumer extends DurableObject {
ws = null;
buf = [];
lastFlush = 0;
cursor = null;
flushing = false;
// Arm the reconnect watchdog first, then connect (idempotent).
async start() {
await this.ctx.storage.setAlarm(Date.now() + RECONNECT_MS);
await this.ensureConnected();
return { connected: this.ws !== null };
}
async ensureConnected() {
if (this.ws) return;
this.cursor ??= (await this.ctx.storage.get("cursor")) ?? null;
const params = new URLSearchParams();
for (const c of WANTED) params.append("wantedCollections", c);
if (this.cursor) params.set("cursor", String(this.cursor));
const resp = await fetch(`${JETSTREAM_URL}?${params}`, {
headers: { Upgrade: "websocket" },
});
const ws = resp.webSocket;
if (!ws) throw new Error(`Jetstream handshake failed: ${resp.status}`);
ws.accept();
this.ws = ws;
ws.addEventListener("message", (e) => this.onMessage(e));
ws.addEventListener("close", () => (this.ws = null));
ws.addEventListener("error", () => (this.ws = null));
}
onMessage(e) {
let ev;
try {
ev = JSON.parse(e.data);
} catch {
return;
}
if (typeof ev.time_us === "number") this.cursor = ev.time_us;
const row = toRow(ev);
if (row) this.buf.push(row);
if (
this.buf.length >= FLUSH_MAX ||
Date.now() - this.lastFlush >= FLUSH_MS
) {
void this.flush();
}
}
// Serialize sends: flush one batch at a time, advancing the cursor on success.
async flush() {
if (this.flushing) return;
this.flushing = true;
try {
while (this.buf.length > 0) {
this.lastFlush = Date.now();
const batch = this.buf.splice(0, this.buf.length);
const batchCursor = this.cursor;
try {
await this.env.BSKY_STREAM.send(batch);
await this.ctx.storage.put("cursor", batchCursor);
} catch (err) {
this.buf.unshift(...batch);
console.error("send failed, will retry", err);
return;
}
}
} finally {
this.flushing = false;
}
}
// Watchdog: reconnect if dropped, flush stragglers, always reschedule.
async alarm() {
try {
await this.ensureConnected();
await this.flush();
} catch (err) {
console.error("alarm error", err);
} finally {
await this.ctx.storage.setAlarm(Date.now() + RECONNECT_MS);
}
}
}
export default {
async fetch(_req, env) {
const stub = env.JETSTREAM.get(env.JETSTREAM.idFromName("singleton"));
return Response.json(await stub.start());
},
async scheduled(_event, env) {
const stub = env.JETSTREAM.get(env.JETSTREAM.idFromName("singleton"));
await stub.start();
},
};import { DurableObject } from "cloudflare:workers";
import type { Pipeline } from "cloudflare:pipelines";
interface Env {
BSKY_STREAM: Pipeline;
JETSTREAM: DurableObjectNamespace<JetstreamConsumer>;
}
// Jetstream collection -> our short event_type. Only these are kept.
const COLLECTION_TO_TYPE: Record<string, string> = {
"app.bsky.feed.post": "post",
"app.bsky.feed.like": "like",
"app.bsky.feed.repost": "repost",
"app.bsky.graph.follow": "follow",
"app.bsky.graph.block": "block",
};
const WANTED = Object.keys(COLLECTION_TO_TYPE);
const JETSTREAM_URL = "https://jetstream2.us-east.bsky.network/subscribe";
const FLUSH_MAX = 500; // rows per send()
const FLUSH_MS = 1000; // flush at least once per second
const RECONNECT_MS = 15000;
// Flatten one Jetstream message into a unified stream row, or null to skip.
function toRow(ev: any) {
if (ev?.kind !== "commit" || !ev.commit) return null;
const c = ev.commit;
const event_type = COLLECTION_TO_TYPE[c.collection];
if (!event_type) return null;
const r = c.record ?? {};
const subject = r.subject;
return {
event_id: `${ev.did}/${c.collection}/${c.rkey}`,
event_type,
did: ev.did ?? null,
operation: c.operation ?? null,
event_time:
typeof ev.time_us === "number"
? new Date(ev.time_us / 1000).toISOString()
: null,
created_at: typeof r.createdAt === "string" ? r.createdAt : null,
text: event_type === "post" && typeof r.text === "string" ? r.text : null,
langs:
event_type === "post" && Array.isArray(r.langs)
? r.langs.join(",")
: null,
subject_uri: typeof subject === "object" ? (subject?.uri ?? null) : null,
subject_did: typeof subject === "string" ? subject : null,
};
}
export class JetstreamConsumer extends DurableObject<Env> {
private ws: WebSocket | null = null;
private buf: Record<string, unknown>[] = [];
private lastFlush = 0;
private cursor: number | null = null;
private flushing = false;
// Arm the reconnect watchdog first, then connect (idempotent).
async start() {
await this.ctx.storage.setAlarm(Date.now() + RECONNECT_MS);
await this.ensureConnected();
return { connected: this.ws !== null };
}
private async ensureConnected() {
if (this.ws) return;
this.cursor ??= (await this.ctx.storage.get<number>("cursor")) ?? null;
const params = new URLSearchParams();
for (const c of WANTED) params.append("wantedCollections", c);
if (this.cursor) params.set("cursor", String(this.cursor));
const resp = await fetch(`${JETSTREAM_URL}?${params}`, {
headers: { Upgrade: "websocket" },
});
const ws = resp.webSocket;
if (!ws) throw new Error(`Jetstream handshake failed: ${resp.status}`);
ws.accept();
this.ws = ws;
ws.addEventListener("message", (e) => this.onMessage(e));
ws.addEventListener("close", () => (this.ws = null));
ws.addEventListener("error", () => (this.ws = null));
}
private onMessage(e: MessageEvent) {
let ev: any;
try {
ev = JSON.parse(e.data as string);
} catch {
return;
}
if (typeof ev.time_us === "number") this.cursor = ev.time_us;
const row = toRow(ev);
if (row) this.buf.push(row);
if (
this.buf.length >= FLUSH_MAX ||
Date.now() - this.lastFlush >= FLUSH_MS
) {
void this.flush();
}
}
// Serialize sends: flush one batch at a time, advancing the cursor on success.
private async flush() {
if (this.flushing) return;
this.flushing = true;
try {
while (this.buf.length > 0) {
this.lastFlush = Date.now();
const batch = this.buf.splice(0, this.buf.length);
const batchCursor = this.cursor;
try {
await this.env.BSKY_STREAM.send(batch);
await this.ctx.storage.put("cursor", batchCursor);
} catch (err) {
this.buf.unshift(...batch);
console.error("send failed, will retry", err);
return;
}
}
} finally {
this.flushing = false;
}
}
// Watchdog: reconnect if dropped, flush stragglers, always reschedule.
async alarm() {
try {
await this.ensureConnected();
await this.flush();
} catch (err) {
console.error("alarm error", err);
} finally {
await this.ctx.storage.setAlarm(Date.now() + RECONNECT_MS);
}
}
}
export default {
async fetch(_req, env): Promise<Response> {
const stub = env.JETSTREAM.get(env.JETSTREAM.idFromName("singleton"));
return Response.json(await stub.start());
},
async scheduled(_event, env): Promise<void> {
const stub = env.JETSTREAM.get(env.JETSTREAM.idFromName("singleton"));
await stub.start();
},
} satisfies ExportedHandler<Env>;为绑定生成类型:
npx wrangler typesyarn wrangler typespnpm wrangler types部署 Worker:
npx wrangler deployyarn wrangler deploypnpm wrangler deploy打开 Worker URL 一次以启动 firehose。cron 触发器使其保持运行:
curl https://bluesky-pipeline.YOUR_SUBDOMAIN.workers.dev命令返回:
{ "connected": true }跟踪日志以观察运行情况:
npx wrangler tailyarn wrangler tailpnpm wrangler tail在第一批事件到达后几分钟内,第一批数据就会落地,pipeline 在此期间预热。
设置 R2 SQL 令牌,然后查询每个表。将 YOUR_WAREHOUSE_NAME 替换为步骤 3 中记下的 warehouse 名称。
export WRANGLER_R2_SQL_AUTH_TOKEN=YOUR_CATALOG_TOKEN
npx wrangler r2 sql query "YOUR_WAREHOUSE_NAME" \
"SELECT COUNT(*) FROM bluesky.bsky_like"
npx wrangler r2 sql query "YOUR_WAREHOUSE_NAME" \
"SELECT text, langs FROM bluesky.bsky_post WHERE langs LIKE '%en%' LIMIT 10"每个表仅包含其事件类型,投影到相关列。单个 pipeline 完成了所有路由。
您使用 Durable Object 消费了高速公共 WebSocket firehose,将其摄取到一个 Pipelines stream,并使用包含多条 SQL 语句的单个 pipeline 将 stream 扇出到五个按事件类型分类的 Iceberg 表。
这种单 stream 到多表的模式可推广到任何带标签的事件源:点击流(按 event_type)、日志(按 service 或 status)或 IoT 遥测(按 device_class)。要扩展它,添加一个 sink 和匹配的 INSERT ... WHERE 语句。
要了解此处使用的 SQL 的更多信息,请参阅 SELECT 语句和管理 pipeline。