Skip to content

非公式本サイトは非公式の日本語ドキュメントであり、Cloudflare 公式サイトではありません。最新情報はdevelopers.cloudflare.comをご確認ください。

ストリームを複数の Iceberg テーブルへ振り分ける

Durable Object で Bluesky Jetstream の firehose を受信し、1 つの Pipelines ストリームに取り込み、複数の SQL 文を持つ 1 つのパイプラインでイベント種別ごとに別の R2 Data Catalog テーブルへ振り分けます。

最終更新 Markdown で表示Agent セットアップ

この例では、公開されている Bluesky Jetstream の firehose(ネットワーク上の投稿、いいね、リポスト、フォロー、ブロックをすべて流すライブ WebSocket ストリーム)を受信し、照会可能な Apache Iceberg テーブルとして R2 Data Catalog に保存します。

Pipelines の基本パターンを学びます。すべてのイベントを 1 つのストリームへ送り、複数の SQL 文を持つ 1 つのパイプラインで、そのストリームをイベント種別ごとの宛先テーブルへ振り分けます(fan out)。種別ごとに別のパイプラインを動かす必要はありません。

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 のバージョンマネージャー

権限の問題を避け、Node.js のバージョンを切り替えられるよう、Voltanvm などの Node バージョンマネージャーを使います。このガイドの後半で説明する Wrangler には、Node バージョン 16.17.0 以降が必要です。

あわせて、R2 Data Catalog と R2 SQL へのアクセスを含む Admin Read & Write 権限の R2 API トークン が必要です。このトークンを各シンクに渡します。Jetstream は公開で認証不要のため、Bluesky アカウントや API キーは不要です。

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. ストリームのスキーマを定義する

ストリームのスキーマは 1 つです。すべてのイベント種別が同じストリームを通るため、スキーマは各種別で使いたいフィールドの和集合になります。あとの振り分け文では、宛先に必要な列だけを選びます。

プロジェクトのルートに 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 を有効にする

シンクは R2 Data Catalog の Iceberg テーブルへ書き込むため、カタログを有効にしたバケットが必要です。

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. ストリーム、シンク、パイプラインを作成する

まず、スキーマファイルからストリームを作成します。

npx wrangler pipelines streams create bsky_events_stream --schema-file schema.json

出力に含まれる stream ID を控えます。次のステップで Worker バインディングを設定するときに使います。

次に、宛先テーブルごとに シンク を 1 つ作成します。各シンクは 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

続いて、ルートごとに 1 つの INSERT 文を含む SQL で、パイプラインを 1 つ作成します。各文は event_type でストリームを絞り込み、そのテーブルに必要な列だけを投影します。

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';

ファイルからパイプラインを作成します。

npx wrangler pipelines create bsky_pipeline --sql-file fanout.sql

1 つのパイプラインが 5 つのテーブルへ書き込みます。あとから新しいイベント種別を足すには、シンクと INSERT 文を 1 つずつ追加します。パイプラインの SQL は作成後に変更できないため、変更する場合はパイプラインを削除して作り直します。詳細は 1 つのストリームを複数テーブルへ振り分ける を参照してください。

5. ストリームを Worker にバインドする

ストリームバインディング、WebSocket 接続を保持する Durable Object、コンシューマーを生かしておく cron トリガー を追加します。<STREAM_ID> をステップ 4 のストリーム ID に置き換えます。

{
  "$schema": "./node_modules/wrangler/config-schema.json",
  "name": "bluesky-pipeline",
  "main": "src/index.ts",
  // Set this to today's date
  "compatibility_date": "2026-09-20",
  "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-09-20"

[[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 を受信する

長時間接続の WebSocket には Durable Object が適しています。ソケットが開いているあいだ常駐し、接続が切れた場合は alarm で再接続します。受信イベントをバッファし、リクエストあたり 5 MB の上限を超えないよう、バッチでストリームへ send() します。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 でテーブルを照会する

最初のデータが届くのは、最初のイベント到着から数分後です。パイプラインのウォームアップに時間がかかります。

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"

各テーブルには、そのイベント種別だけが、必要な列に投影されて入ります。振り分けは 1 つのパイプラインがすべて行います。

まとめ

高頻度の公開 WebSocket firehose を Durable Object で受信し、1 つの Pipelines ストリームに取り込み、複数の SQL 文を持つ 1 つのパイプラインで、イベント種別ごとに 5 つの Iceberg テーブルへ振り分けました。

この 1 ストリーム対複数テーブルのパターンは、タグ付きのイベントソース全般に使えます。クリックストリーム(event_type ごと)、ログ(servicestatus ごと)、IoT テレメトリ(device_class ごと)などです。拡張するには、シンクと対応する INSERT ... WHERE 文を追加します。

ここで使った SQL の詳細は、SELECT 文パイプラインの管理 を参照してください。

役に立ちましたか?