Cloudflare K2:无服务器事件流

传统远程过程调用(RPC)架构有一个核心难题:生产者和消费者必须在规模和时间上保持同步。如果生产者发送的数据量超出消费者的处理能力,或者消费者或下游服务不可用,事件就会被丢弃。当多个消费者需要各自独立处理数据时,这个问题会进一步放大。比如,电商后端在交易完成时会发出事件,分析系统和欺诈检测服务都需要读取这些事件。
解决办法是让生产者和消费者解耦——在两者之间加入一个服务,由它承接写入,同时让各个读取方按自己的节奏独立消费。
今天,我们正式推出 Cloudflare K2 公测版来解决这个问题。K2 是开发者平台上的持久化事件流原语。你把事件发送到 K2 stream,它以有序日志的形式存储事件。消费者可以通过多种方式读取,例如将读取任务分摊到一组消费者上,或者把所有消息投递给每个消费者。它完全 serverless,能扩展到海量数据规模,并支持长期保留,即使消费者长时间停机也不会丢数据。
在底层,K2 基于 R2 对象存储实现了分区、持久化的日志,因此能够扩展到极大的存储规模。
如果你已经跃跃欲试,可以按照这份指南,几秒钟内创建你的第一个 stream。
边缘上的事件流
我们最初构建 K2 是因为自己需要一个边缘上的持久化缓冲区,最初是作为 Basin Pipelines 的数据接入层。Pipelines 由一个基于拉取模式的流处理引擎驱动,这意味着在事件被读取、转换并写入 R2 之前,需要有另一个系统来存储它们。而且由于我们承诺事件一旦进入 Pipelines Stream 就绝不丢弃,这个存储层必须是持久化的——也就是说不能丢数据——并且要能支撑可能相当长的时间。
在大多数公司,这里通常会部署 Apache Kafka。但 Pipelines 运行在 Cloudflare 的边缘网络上,该网络横跨 335 多个城市中的大量服务器。由于架构独特,我们通常无法运行像 Kafka 这样的传统分布式系统软件,必须重新思考这些系统的构建和运营方式。
具体来说,对于有状态服务,Cloudflare 的全球基础设施存在一些挑战:我们获得的机器切片相对较小,且这些机器生命周期短暂,网络通信往往基于公共互联网。但我们的基础设施也有几大优势:它在全球各地都靠近用户,并具备极强的水平扩展能力。
在设计最终演变为 K2 的持久化缓冲系统时,我们决定依赖已有的强大状态原语:R2。像 R2 这样的对象存储系统,将极高的数据持久性(11 个 9!)与强一致性 API 相结合。将复制和共识下放至存储层,使得应用层(在本例中为 K2)变得极其简单、廉价且高性能。其次,这种分离计算和存储的方式,使得两者可以独立扩展,从而让我们能以低成本存储海量的历史数据。
如何在对象存储之上构建日志?首要问题是,R2 与其他对象存储一样,不支持追加操作,而这是日志的标准操作。相反,我们必须写入完整的文件或分段,且文件大小需足够大,以抵消每次读写产生的开销。我们的做法是,先在边缘服务中通过内存累积写入数据。在等待短时间内让数据到齐后,我们将所有事件作为一个分段文件写入。利用 R2 的原子操作,我们实现了事件的顺序排序和严格递增的偏移量,无需额外的协调服务。
尽管基于 R2 构建优势众多,但有一个缺点:较高的生产延迟。写入对象存储比写入本地磁盘慢,且我们必须在本地批次数据积累完成后才开始写入。在 K2 的初始版本中,这使得响应时间的第 99 百分位下的生产延迟增加了约 1 秒。
我们将通过即将发布的技术深度剖析分享更多关于 K2 设计的细节。
流、队列还是管道?
### 云端流式新范式:K2 的无服务器事件流Cloudflare 已经拥有多种成熟的异步传输原语,例如 Queues 和 Basin Pipelines。那么,什么场景下应该优先选择 K2?
Queues 和 K2 Streams 在表面功能上有些相似:它们都能接收事件、持久化存储并投递给消费者。Queues 的核心设计理念是追踪那些耗时长或资源昂贵的单个工作项,并支持对它们的异步处理。举个例子,在图像应用中,你可以将用户上传的请求加入队列,交由后端服务进行实际处理。在这种粒度下,它支持复杂的逻辑,例如重试、延迟控制,以及针对失败任务设置的死信队列。
相比之下,K2 侧重于大规模的数据流动、长期保留和扇出消费。消息以批量(Batch)方式生产和消费,这种设计以牺牲单条消息级别的重试能力为代价,换取了极高的处理效率。这种批量传输机制也意味着生产者端的延迟通常高于传统队列。
Basin Pipelines 是一款无服务器摄取服务。你可以向其中发送 Pipeline JSON 事件,这些事件可以被转换后写入 R2 或 Basin Catalog。当最终目标是将事件写入对象存储或 Iceberg 表时,推荐选用 Pipelines;若需自定义处理逻辑或写入其他目标系统,则推荐选用 K2。
快速入门
使用 K2 的第一步是创建一个 Stream(流)。在一个账户下可以创建多个流,以满足不同的使用场景或处理不同类别的事件。你可以使用 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 ${CLOUDFLARE_API_TOKEN}" \
-H "Content-Type: application/json" \
--data '{
"worker_id": "analytics-1",
"max_records": 100
}' {
"result": {
"batch_id": "b7e4c9210a3f468d95c2e1068fdb734a",
"leased_until_ms": 1790633929437,
"records": [
{
"timestamp_ms": 1790633629168,
"content": "eyJldmVudCI6InBhZ2VfdmlldyJ9",
"headers": {
"content-type": "application/json"
}
},
...
]
},
"success": true,
"errors": [],
"messages": []
}客户端调用 consume 时,会获得这批事件为期 5 分钟的租约。客户端可以做的事有三种:
- ack 确认该批次,将其标记为已处理,确保不会再重新投递
- nack(负确认)该消息,表示处理失败,请求重新投递
- 延长其租约(lease),以防处理需要更多时间
$ curl -X POST "https://d78b09ee1f50430e9ec92a8af92b0231.k2.cloudflarestorage.com/subscriptions/ee13f761783d3823a447a47b572ebf76/batches/b7e4c9210a3f468d95c2e1068fdb734a/ack" \
-H "Authorization: Bearer ${CLOUDFLARE_API_TOKEN}" \
-H "Content-Type: application/json" \
--data '{ "worker_id": "analytics-1" }'这是消费 K2 的一种方式:将工作分配给多个消费者,每个消费者获取一部分数据。另一种方式是使用发布/订阅(pub/sub)模式,为每个消费者创建独立的订阅,这样每个消费者都能看到所有消息。你也可以结合这两种方法,建立多个独立的消费者池。
关于 API 的完整细节,请参阅 K2 文档。
定价与可用性
K2 目前处于公开测试阶段(public beta),适用于拥有 Workers 付费订阅的账户,并受以下限制:
- 最大存储用量为 10GB
- 每个流的产生(produce)速率为 30 MB/s
如需提高限制,请在 Discord 上联系团队或填写限制提升申请。
K2 的使用在测试期间不计费。开始计费后,预计定价如下:
定价 | |
数据产生 | $0.04 / GB |
数据消费 | $0.04 / GB |
数据保留 | $0.02 / GB / 月 |
下一步计划
未来几个月,我们为 K2 制定了一份激动人心的路线图,包括:
- 更高的写入并行度,支持每秒多 GB 的流
- 消息键和基于键的顺序保证
- 基于推送的 worker 消费者
- Express 层级,具有更低的产生和端到端延迟
- 对 Apache Kafka 客户端的即插即用支持
我们期待看到您在 K2 上构建的内容!请在 Cloudflare Discord 上分享您的反馈。