Cloudflare K2: serverless event streams
Cloudflare K2: Serverless Event Streams
Cloudflare K2:无服务器事件流
With traditional Remote Procedure Call (RPC) architectures, there exists a core challenge: producers and consumers must align in scale and in time. If your producers send too much data for your consumers to handle or if your consumers or downstream services become unavailable, events are dropped. This problem is compounded with multiple consumers that need to independently process the data. For example, an ecommerce backend may emit events when transactions are completed, which need to be read by an analytics system and a fraud detection service.
在传统的远程过程调用(RPC)架构中,存在一个核心挑战:生产者和消费者必须在规模和时间上保持一致。如果生产者发送的数据量超过了消费者的处理能力,或者消费者及下游服务不可用,事件就会丢失。当多个消费者需要独立处理数据时,这个问题会变得更加复杂。例如,电子商务后端在交易完成时可能会发出事件,这些事件需要被分析系统和欺诈检测服务同时读取。
We can solve this by decoupling our producers and consumers — inserting a service in the middle that absorbs writes while allowing independent readers to consume at their own pace. Today we are launching Cloudflare K2 in public beta to solve this problem. K2 is a durable event streaming primitive on the Developer Platform. You send events to a K2 stream, which stores them as an ordered log. Consumers can read them in a variety of ways, for example by splitting up reads across a set of consumers, or delivering all messages to all consumers. It’s fully serverless, scales to vast quantities of data, and supports long-term retention, so even long periods of consumer downtime do not lose data.
我们可以通过解耦生产者和消费者来解决这个问题——在中间插入一个服务,既能吸收写入的数据,又允许独立的读取者按自己的节奏进行消费。今天,我们正式推出 Cloudflare K2 公测版来解决这一问题。K2 是开发者平台上的一个持久化事件流原语。你将事件发送到 K2 流,它会将其存储为有序日志。消费者可以通过多种方式读取这些数据,例如将读取任务分配给一组消费者,或者将所有消息分发给所有消费者。它是完全无服务器的,能够扩展到海量数据,并支持长期保留,因此即使消费者长时间停机,数据也不会丢失。
Under the hood, K2 implements a partitioned, durable log on top of R2 object storage, which allows it to scale to huge volumes of storage. If you’re ready to get started, you can create your first stream in seconds by following the guide here.
在底层,K2 在 R2 对象存储之上实现了一个分区化的持久日志,这使其能够扩展到巨大的存储容量。如果你准备好开始使用,可以按照此处的指南在几秒钟内创建你的第一个流。
Streams on the edge
边缘流
We first built K2 because we needed a durable buffer on the edge, initially to serve as the ingestion layer for Basin Pipelines. Pipelines is powered by a stream processing engine that operates on a pull-based model, which means some other system has to store events before they are read, transformed, and written to R2. And because we commit to never dropping events once they’re accepted into the Pipelines Stream, that storage has to be durable — meaning it can’t lose data — over potentially long periods of time.
我们最初构建 K2 是因为需要在边缘端建立一个持久化缓冲区,最初是作为 Basin Pipelines 的摄取层。Pipelines 由一个基于拉取模型的流处理引擎驱动,这意味着在事件被读取、转换并写入 R2 之前,必须有其他系统先存储这些事件。由于我们承诺一旦事件被 Pipelines Stream 接收就绝不丢失,因此该存储必须是持久的——意味着在可能很长的时间内不能丢失数据。
This is where most companies would deploy Apache Kafka. However, Pipelines runs on the Cloudflare edge, which spans a huge number of servers across over 335 cities. Our unique architecture means we often cannot run traditional distributed systems software like Kafka, and need to rethink how these systems are built and operated.
大多数公司在这种情况下会部署 Apache Kafka。然而,Pipelines 运行在 Cloudflare 边缘,覆盖了 335 多个城市的庞大服务器集群。我们独特的架构意味着我们通常无法运行像 Kafka 这样的传统分布式系统软件,因此需要重新思考这些系统的构建和运行方式。
For stateful services, in particular, Cloudflare’s global infrastructure presents some challenges: we get relatively small slices of machines, those machines are relatively ephemeral, and networking is often over the public Internet. But our infrastructure also has a few superpowers: it’s close to users wherever they are in the world and has an incredible capacity to scale horizontally.
特别是对于有状态服务,Cloudflare 的全球基础设施带来了一些挑战:我们获得的机器资源相对碎片化,这些机器相对短暂,且网络通常通过公共互联网传输。但我们的基础设施也有一些“超能力”:它在世界各地都靠近用户,并具备极强的水平扩展能力。
In designing the durable buffering system that became K2, we decided to rely on the powerful state primitive we already have: R2. Object storage systems like R2 combine extremely durable storage (11 9s!) with strongly consistent APIs. Offloading replication and consensus to the storage layer allows us to make the application layer (K2 in this case) radically simpler, cheaper, and higher performance. A secondary benefit is that it separates compute and storage, meaning each can be scaled independently. This allows us to store vast quantities of historical data at low cost.
在设计最终成为 K2 的持久化缓冲系统时,我们决定依赖我们已有的强大状态原语:R2。像 R2 这样的对象存储系统结合了极高的持久性(11 个 9!)和强一致性 API。将复制和共识机制卸载到存储层,使我们能够让应用层(本例中为 K2)变得极其简单、廉价且高性能。另一个好处是它实现了计算与存储的分离,意味着两者可以独立扩展。这使我们能够以低成本存储海量的历史数据。
How do we build a log on top of object storage? An immediate issue is that R2 — like other object stores — does not support appends, the standard operation on a log. Instead, we must write complete files, or segments, that are large enough to overcome the cost of writing and reading each one. We do this by first accumulating writes in-memory on an edge service. After waiting a short period for data to arrive, we write all events as a segment file. We achieve ordering and strictly incrementing offsets using R2’s atomic operations without needing a separate coordination service.
我们如何基于对象存储构建日志?一个直接的问题是,R2(像其他对象存储一样)不支持追加(append)操作,而这是日志的标准操作。相反,我们必须写入完整的“段”文件,这些文件必须足够大,以抵消每次读写的开销。我们通过首先在边缘服务的内存中累积写入来实现这一点。在等待一小段时间以接收数据后,我们将所有事件写入为一个段文件。我们利用 R2 的原子操作实现了排序和严格递增的偏移量,而无需额外的协调服务。
While building on R2 has many advantages, there is one downside: higher produce latencies. Writing to object storage is slower than a local disk, and we have to wait for the local batch to accumulate before starting the write. In our initial release of K2, this adds up to about 1 second of produce latency at the 99th percentile of response times.
虽然基于 R2 构建有很多优势,但也有一个缺点:更高的生产延迟。写入对象存储比本地磁盘慢,而且我们必须等待本地批次累积完成后才能开始写入。在 K2 的初始版本中,这导致 99 分位响应时间的生产延迟增加了约 1 秒。
We will be sharing more details on the design of K2 in an upcoming technical deep dive.
我们将在即将发布的技术深度解析中分享更多关于 K2 设计的细节。
Streams, Queues, or Pipelines?
流、队列还是管道?
Cloudflare has several existing asynchronous delivery primitives, including Queues and Basin Pipelines. When should you reach for K2 instead of these existing products?
Cloudflare 已经有几种现有的异步交付原语,包括 Queues 和 Basin Pipelines。什么时候应该选择 K2 而不是这些现有产品呢?
There are some superficial similarities between Queues and K2 Streams: both receive events, durably store them, and deliver them to consumers. Queues are designed around tracking individual items of expensive or time consuming work that need to be asynchronously completed. For example, an image processing application may enqueue a user request to be handled by the actual image processing service. They support complex logic on the grain of a particular work item, like retries, delays, and dead-letter queues for failed attempts.
Queues 和 K2 Streams 在表面上有一些相似之处:它们都接收事件、持久化存储并交付给消费者。Queues 的设计初衷是跟踪需要异步完成的昂贵或耗时的单个工作项。例如,图像处理应用程序可以将用户请求排队,由实际的图像处理服务进行处理。它们支持针对特定工作项的复杂逻辑,如重试、延迟和失败尝试的死信队列。
K2, by contrast, is designed for high-scale data movement, long-term retention, and fan-out consumption. Messages are produced and consumed as batches — enabling efficient processing at the expense of message-level retries. This batching also drives higher producer latency than for queues.
相比之下,K2 专为大规模数据传输、长期保留和扇出(fan-out)消费而设计。消息以批次形式生产和消费——这实现了高效处理,但代价是牺牲了消息级的重试机制。这种批处理方式也导致了比队列更高的生产延迟。
Basin Pipelines is a serverless ingestion service. You can send your Pipeline JSON events, which can be transformed and written to R2 or a Basin Catalog. We recommend Pipelines when the end result is writing your events to object storage or Iceberg tables, and K2 when doing custom processing or writing to other destinations.
Basin Pipelines 是一项无服务器摄取服务。你可以发送 Pipeline JSON 事件,这些事件可以被转换并写入 R2 或 Basin Catalog。如果最终目标是将事件写入对象存储或 Iceberg 表,我们推荐使用 Pipelines;如果需要进行自定义处理或写入其他目标,则推荐使用 K2。
Getting started
开始使用
Using K2 involves first creating a stream. You can have many streams across your account for different use cases or types of events. Streams can be created via cf, Wrangler, the dashboard, or API. Let’s take the example of collecting and processing product analytics. First, we’ll create a stream with…
使用 K2 首先需要创建一个流。你可以在账户中创建多个流,以应对不同的用例或事件类型。流可以通过 cf、Wrangler、仪表板或 API 创建。让我们以收集和处理产品分析数据为例。首先,我们将创建一个流,其中…