跳转到内容
搜索文档

将 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 版本管理器

使用 Voltanvm 等 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)、日志(按 servicestatus)或 IoT 遥测(按 device_class)。要扩展它,添加一个 sink 和匹配的 INSERT ... WHERE 语句。

要了解此处使用的 SQL 的更多信息,请参阅 SELECT 语句管理 pipeline

这篇文档对您有帮助吗?