この例では、公開されている 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]
- Cloudflare アカウント ↗ に登録します。
Node.js↗ をインストールします。
Node.js のバージョンマネージャー
権限の問題を避け、Node.js のバージョンを切り替えられるよう、Volta ↗ や nvm ↗ などの Node バージョンマネージャーを使います。このガイドの後半で説明する Wrangler には、Node バージョン 16.17.0 以降が必要です。
あわせて、R2 Data Catalog と R2 SQL へのアクセスを含む Admin Read & Write 権限の R2 API トークン が必要です。このトークンを各シンクに渡します。Jetstream は公開で認証不要のため、Bluesky アカウントや API キーは不要です。
次のコマンドで、新しい 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@latestストリームのスキーマは 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 }
]
}シンクは R2 Data Catalog の Iceberg テーブルへ書き込むため、カタログを有効にしたバケットが必要です。
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 でデータを照会するときに使います。
まず、スキーマファイルからストリームを作成します。
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 バインディングを設定するときに使います。
次に、宛先テーブルごとに シンク を 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.sqlyarn wrangler pipelines create bsky_pipeline --sql-file fanout.sqlpnpm wrangler pipelines create bsky_pipeline --sql-file fanout.sql1 つのパイプラインが 5 つのテーブルへ書き込みます。あとから新しいイベント種別を足すには、シンクと INSERT 文を 1 つずつ追加します。パイプラインの SQL は作成後に変更できないため、変更する場合はパイプラインを削除して作り直します。詳細は 1 つのストリームを複数テーブルへ振り分ける を参照してください。
ストリームバインディング、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 * * * *"]長時間接続の WebSocket には Durable Object が適しています。ソケットが開いているあいだ常駐し、接続が切れた場合は alarm で再接続します。受信イベントをバッファし、リクエストあたり 5 MB の上限を超えないよう、バッチでストリームへ send() します。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 typesWorker をデプロイします。
npx wrangler deployyarn wrangler deploypnpm wrangler deployWorker の URL を一度開き、firehose を起動します。以降は cron トリガーが稼働を維持します。
curl https://bluesky-pipeline.YOUR_SUBDOMAIN.workers.devコマンドは次を返します。
{ "connected": true }ログを追跡して動作を確認します。
npx wrangler tailyarn wrangler tailpnpm wrangler tail最初のデータが届くのは、最初のイベント到着から数分後です。パイプラインのウォームアップに時間がかかります。
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 ごと)、ログ(service や status ごと)、IoT テレメトリ(device_class ごと)などです。拡張するには、シンクと対応する INSERT ... WHERE 文を追加します。