インフラ・DevOps
Cloudflare K2: サーバーレス イベントストリーム
Cloudflare K2: serverless event streams (blog.cloudflare.com)
要約
Cloudflareは、従来のRPCアーキテクチャにおけるプロデューサーとコンシューマーのスケーリングと同期の課題を解決するため、サーバーレスイベントストリーミングサービス「K2」をパブリックベータでローンチしました。K2は、イベントを順序付けられたログとして耐久性高く保存し、コンシューマーが自身のペースでデータを読み取れるようにします。
全文翻訳
従来のRPC(Remote Procedure Call)アーキテクチャには、プロデューサーとコンシューマーが規模と時間において一致しなければならないという根本的な課題が存在します。プロデューサーがコンシューマーの処理能力を超えたデータを送信したり、コンシューマーや下流のサービスが利用できなくなったりすると、イベントは失われます。これは、データを独立して処理する必要がある複数のコンシューマーがいる場合に、さらに複雑になります。例えば、eコマースのバックエンドでは、トランザクションが完了した際にイベントを発行する必要があり、これを分析システムと不正検出サービスが読み取る必要があります。
この問題は、プロデューサーとコンシューマーを疎結合にすることで解決できます。その間にサービスを挿入し、書き込みを吸収しながら、独立したリーダーが自身のペースで消費できるようにします。
本日、この問題解決のためにCloudflare K2をパブリックベータでローンチします。K2は、開発者プラットフォーム上の耐久性のあるイベントストリーミングプリミティブです。K2ストリームにイベントを送信すると、順序付けられたログとして保存されます。コンシューマーは、例えば一連のコンシューマー間で読み取りを分割したり、すべてのメッセージをすべてのコンシューマーに配信したりするなど、さまざまな方法で読み取ることができます。完全にサーバーレスで、膨大なデータ量にスケールし、長期保存をサポートするため、コンシューマーが長期間ダウンしてもデータが失われることはありません。
内部的には、K2はR2オブジェクトストレージの上にパーティション化された耐久性のあるログを実装しており、これにより巨大なストレージ容量にスケールできます。
準備ができた方は、こちらのガイドに従って数秒で最初のストリームを作成できます。
エッジでのストリーム
K2を最初に構築したのは、エッジに耐久性のあるバッファーが必要だったからです。当初はBasin Pipelinesの取り込みレイヤーとして機能していました。Pipelinesはプルベースのモデルで動作するストリーム処理エンジンによって駆動されているため、何らかのシステムがイベントを読み取り、変換し、R2に書き込む前にイベントを保存する必要があります。そして、Pipelinesストリームに取り込まれたイベントを決して失わないことを約束しているため、そのストレージは、長期間にわたって耐久性がある(つまり、データを失わない)必要があります。
ここでほとんどの企業はApache Kafkaをデプロイするでしょう。しかし、PipelinesはCloudflareのエッジで実行されており、335都市以上に広がる膨大な数のサーバーにまたがっています。私たちのユニークなアーキテクチャは、Kafkaのような従来の分散システムソフトウェアを実行できないことが多く、これらのシステムがどのように構築され、運用されるかを再考する必要があることを意味します。
特にステートフルなサービスにとって、Cloudflareのグローバルインフラストラクチャはいくつかの課題をもたらします。比較的少量のマシンリソースしか得られず、それらのマシンは比較的一時的であり、ネットワークはパブリックインターネット経由であることがよくあります。しかし、私たちのインフラストラクチャにはいくつかのスーパーパワーもあります。世界中のどこにいてもユーザーに近いこと、そして水平方向にスケールする驚異的な能力を持っていることです。
K2となった耐久性のあるバッファリングシステムを設計するにあたり、私たちはすでに持っている強力なステートプリミティブであるR2に依存することにしました。R2のようなオブジェクトストレージシステムは、非常に耐久性の高いストレージ(9の後に11個の9!)と強く一貫性のあるAPIを組み合わせています。レプリケーションとコンセンサスをストレージレイヤーにオフロードすることで、アプリケーションレイヤー(この場合はK2)を根本的にシンプル、安価、そして高性能にすることができます。二次的な利点は、コンピュートとストレージを分離しているため、それぞれを独立してスケールできることです。これにより、大量の履歴データを低コストで保存できます。
オブジェクトストレージの上にログを構築するにはどうすればよいでしょうか?すぐに問題となるのは、R2(他のオブジェクトストアと同様)が、ログの標準的な操作であるアペンドをサポートしていないことです。代わりに、各セグメントの書き込みと読み取りのコストを克服するのに十分な大きさのファイル、またはセグメントを書き込む必要があります。これを実現するために、まずエッジサービスでメモリ内に書き込みを蓄積します。データが到着するまで短い時間待った後、すべてのイベントをセグメントファイルとして書き込みます。個別のコーディネーションサービスを必要とせずに、R2のアトミック操作を使用して順序と厳密に増加するオフセットを実現します。
R2上に構築することには多くの利点がありますが、1つの欠点があります。それは、プロデュースレイテンシが高くなることです。オブジェクトストレージへの書き込みはローカルディスクよりも遅く、書き込みを開始する前にローカルバッチが蓄積されるのを待つ必要があります。K2の初期リリースでは、これは応答時間の99パーセンタイルで約1秒のプロデュースレイテンシにつながります。
K2の設計に関する詳細は、今後の技術的な詳細解説で共有します。
ストリーム、キュー、またはパイプライン?
Cloudflareには、QueuesやBasin Pipelinesを含むいくつかの既存の非同期配信プリミティブがあります。これらの既存の製品の代わりにK2を使用するのはどのような場合でしょうか?
QueuesとK2 Streamsの間には、表面的な類似点がいくつかあります。どちらもイベントを受け取り、耐久性高く保存し、コンシューマーに配信します。Queuesは、非同期に完了する必要がある高コストまたは時間のかかる作業の個々のアイテムを追跡することを中心に設計されています。例えば、画像処理アプリケーションは、実際の画像処理サービスによって処理されるユーザーリクエストをキューに入れる場合があります。これらは、リトライ、遅延、失敗した試行のためのデッドレターキューなど、特定の作業アイテムの粒度で複雑なロジックをサポートします。
対照的に、K2は、高スケールのデータ移動、長期保存、およびファンアウト消費のために設計されています。メッセージはバッチとして生成および消費され、メッセージレベルのリトライを犠牲にして効率的な処理を可能にします。このバッチ処理は、キューよりも高いプロデューサーレイテンシも引き起こします。
Basin Pipelinesは、サーバーレスの取り込みサービスです。Pipeline JSONイベントを送信でき、これらは変換されてR2またはBasin Catalogに書き込まれます。最終的な結果がオブジェクトストレージまたはIcebergテーブルへのイベントの書き込みである場合はPipelinesを推奨し、カスタム処理や他の宛先への書き込みを行う場合はK2を推奨します。
はじめに
K2を使用するには、まずストリームを作成します。アカウント全体で、さまざまなユースケースやイベントの種類に対して多数のストリームを持つことができます。ストリームは、cf、Wrangler、ダッシュボード、またはAPIを介して作成できます。
製品分析の収集と処理を例にとってみましょう。まず、cfでストリームを作成します。
$ cf k2 streams create --name app_events --http-enabled
{
"id": "d78b09ee1f50430e9ec92a8af92b0231",
"name": "app_events",
"retention_seconds": 604800,
"endpoint": "https://d78b09ee1f50430e9ec92a8af92b0231.k2.cloudflarestorage.com",
"http": {
"enabled": true,
"authentication": false
},
"worker_binding": {
"enabled": true
},
"created_at": "2026-09-28T15:14:39.053Z",
"modified_at": "2026-09-28T15:14:39.053Z"
}
ストリームが作成されたら、HTTP APIまたはWorkerバインディングを介してそれにプロデュースを開始できます。例えば、Workerバインディングを設定し、次のようにプロデュースできます。
const result = await env.EVENTS.send([
{
content: new TextEncoder().encode(
JSON.stringify({
event: "page_view",
path: new URL(request.url).pathname,
timestamp: Date.now(),
})
),
headers: {
"content-type": "application/json",
},
},
]);
if (!result.success) {
console.error(`Produce failed: ${result.error.message}`);
return new Response("Failed to record event", {
status: result.error.retryable ? 503 : 500,
});
}
K2はデータをバイトとして表現するため、アプリケーションに適した任意の形式またはエンコーディングを使用できます。
イベントがストリームに入ったら、サブスクリプションを作成できます。サブスクリプションは、コンシューマー間で作業を分割し、読み取り並列処理を可能にします。これにより、単一サーバーで処理できる以上の負荷を処理するために、複数のリーダーにスケールアウトできます。
HTTP APIを介してサブスクリプションを作成できます。
$ curl -X POST "https://d78b09ee1f50430e9ec92a8af92b0231.k2.cloudflarestorage.com/subscriptions" \
-H "Authorization: Bearer ${CLOUDFLARE_API_TOKEN}" \
-H "Content-Type: application/json" \
--data '{ "name": "analytics_processor", "start_at": { "type": "earliest" } }'
{
"result": {
"id": "ee13f761783d3823a447a47b572ebf76"
},
"success": true,
"errors": [],
"messages": []
}
サブスクリプションが作成されたら、各コンシューマーからそれをポーリングできます。
$ curl -X POST "https://4d8f5394e3e733debdeca9c65c5b7439.k2.cloudflarestorage.com/subscriptions/ee13f761783d3823a447a47b572ebf76/consume" \
-H "Authorization: Bearer ${CLO