跳转到内容
搜索文档

将 stream 扇出到多个 Iceberg 表

在 Durable Object 中消费 Bluesky Jetstream firehose,将其摄取到单个 Pipelines stream,并使用包含多条 SQL 语句的一个 pipeline 按类型将事件路由到独立的 R2 Data Catalog 表。

最后更新 查看 MarkdownAgent 设置

在本示例中,您将消费公共 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]

前提条件

  1. 注册 Cloudflare 账户 ↗。
  2. 安装 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 是公开且无需身份验证的。

1. 创建新的 Worker 项目

通过运行以下命令创建新的 Worker 项目:

npm 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@latest

2. 定义 stream schema

Stream 只有一个 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 }
	]
}

3. 创建 R2 存储桶并启用 R2 Data Catalog

您的 sink 将数据写入 R2 Data Catalog 中的 Iceberg 表,因此您需要一个启用了 catalog 的存储桶。

创建名为 bluesky-pipeline 的 R2 存储桶:

npx wrangler r2 bucket create bluesky-pipeline

在存储桶上启用 R2 Data Catalog:

npx wrangler r2 bucket catalog enable bluesky-pipeline

运行此命令时,记下 Warehouse name。您将需要它来使用 R2 SQL 查询数据。

4. 创建 stream、sink 和 pipeline

首先,从 schema 文件创建 stream:

npx 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.sql

一个 pipeline 写入五个表。要添加新的事件类型,添加一个 sink 和一条匹配的 INSERT 语句。Pipeline SQL 创建后无法修改,因此您需要删除并重新创建 pipeline 来更改它。有关更多信息,请参阅将一个 stream 路由到多个表。

5. 将 stream 绑定到 Worker

添加 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 * * * *"]

6. 在 Durable Object 中消费 firehose

Durable Object 是长期 WebSocket 的理想宿主。它在 socket 打开时保持驻留,alarm 在连接断开时重新连接。缓冲传入事件并批量 send() 到 stream,以保持在每次请求 5 MB 的限制以下。持久化 Jetstream time_us 游标,以便重连时无间隙恢复。

将 src/index.ts 的内容替换为以下内容:

src/index.jsjs
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();
	},
};
src/index.tsts
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 types

7. 部署并启动消费者

部署 Worker:

npx wrangler deploy

打开 Worker URL 一次以启动 firehose。cron 触发器使其保持运行:

curl https://bluesky-pipeline.YOUR_SUBDOMAIN.workers.dev

命令返回:

{ "connected": true }

跟踪日志以观察运行情况:

npx wrangler tail

8. 使用 R2 SQL 查询表

在第一批事件到达后几分钟内,第一批数据就会落地,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。

这篇文档对您有帮助吗?