Cloudflare K2: how to adopt serverless event streams
Explore with AI
Cloudflare K2 is a durable, serverless event stream that Cloudflare announced in public beta on 1 October 2026: producers write batches of records to an ordered log stored on R2, and up to 100 subscriptions read it at their own pace with at-least-once delivery. Use it for high-volume data movement, fan-out to several consumers and replay within a retention window of up to 30 days, and keep Cloudflare Queues for per-message work that needs retries, delays and dead-letter queues. It needs a Workers Paid plan, beta accounts are capped at 10 GB of storage and 30 MB/s of produce per stream, and produce latency is about 1 second at p99.
On this page
- What Cloudflare K2 is: an ordered log your consumers read at their own pace
- K2, Queues or Pipelines: picking the right Cloudflare primitive
- One stream, many readers: shared and fan-out subscriptions
- Standing up a stream, a producer and a consumer
- The delivery rules your consumers have to handle
- Beta limits and the pricing Cloudflare expects
- What goes wrong when teams first adopt K2
- Watching a K2 pipeline once it carries production traffic
- What is on the K2 roadmap
- Common questions
Cloudflare K2 is a serverless event streaming service that Cloudflare announced in public beta on 1 October 2026. You send batches of records to a stream. K2 stores them as an ordered, durable log on R2 object storage, and consumers read that log through subscriptions, each at its own pace. It is available now to accounts on the Workers Paid plan, and usage is not billed during the beta.
K2 suits teams that need one write to reach several readers, that need to absorb bursts their consumers can’t keep up with, or that need to replay events after an outage. This guide covers what K2 is, how it compares with Cloudflare Queues, Pipelines and Kafka, how to set up a stream and a working consumer today, and the limits and delivery rules that will shape your design.
What Cloudflare K2 is: an ordered log your consumers read at their own pace
K2 has four ideas, and once you have them the API is easy to follow.
- Stream. An ordered, durable, append-only log of records. Each stream has a name, a retention period and an ID that K2 assigns when you create it. Producers and consumers use the ID to reach the stream.
- Record. A binary
contentpayload plus optional stringheaders. K2 stores bytes, so you choose the encoding: JSON, Protobuf, anything. Headers describe the content without anyone having to deserialise it, for example the encoding, an event type or a routing key. - Subscription. A read position in the stream. Every subscription receives every record. Consumers that share one subscription split the records between them.
- Lease. When a consumer asks for a batch, K2 leases it to that consumer for 5 minutes. The consumer acks the batch when it is done, nacks it if processing failed, or extends the lease if it needs more time.
The main difference from a queue is what reading does. In a queue, consuming a message removes it. In K2, writes append to the log and reads only move a pointer forward. That gives you three properties: a slow reader never blocks writes, many readers can consume the whole stream without affecting each other, and you can replay anything still inside the retention window.
Why K2 sits on R2, and the latency that costs
Cloudflare first built K2 as the durable buffer in front of Basin Pipelines. Most companies would use Apache Kafka for that job, but Cloudflare’s edge runs across more than 335 cities on small, short-lived slices of machines, and Kafka doesn’t fit there. So K2 hands replication and consensus to R2, which offers 11 nines of durability and strongly consistent APIs. That keeps compute and storage separate, so each scales on its own.
R2, like other object stores, has no append operation. K2 collects writes in memory on an edge service for a short time, writes them out as a segment file, and uses R2’s atomic operations to keep order and strictly increasing offsets without a separate coordination service. The cost is produce latency. At launch, a produce call takes about 1 second at the 99th percentile.
K2, Queues or Pipelines: picking the right Cloudflare primitive
Cloudflare has several asynchronous delivery products, and Cloudflare’s announcement says when to use each.
- Use K2 for high-volume data movement, long retention and fan-out to several consumers, and when you write your own processing or send data somewhere other than object storage. Records are produced and consumed in batches, which keeps processing efficient.
- Use Queues when each message is a unit of expensive or slow work that someone has to finish, such as an image to process. Queues track individual items and support retries, delays and dead-letter queues for each item. K2 trades away message-level retries: you ack or nack a whole batch.
- Use Pipelines when the end result is your events written to R2 or to Iceberg tables. Pipelines accepts JSON events, transforms them and writes them out for you.
If you run Kafka today, check the roadmap before you plan a migration. Drop-in support for Kafka clients, message keys with key-based ordering, and multi-GB/s streams are all listed as coming in the next few months. Right now you produce over HTTP or a Workers binding, and you consume by polling an HTTP endpoint.
One stream, many readers: shared and fan-out subscriptions
Take the example from the announcement: an ecommerce backend emits an event for every completed transaction, and both an analytics system and a fraud detection service need every event.
- Create one stream,
transactions. - Create a subscription called
analyticsand run four consumer processes against it, each with its ownworker_id. They share the work, and each sees part of the data. - Create a second subscription called
fraud-detectionwith its own pool of consumers. It receives every record too, and it reads at its own speed.
If fraud detection goes down for a day, analytics keeps going and fraud detection catches up from where it stopped when it comes back, as long as the records are still inside the retention window. You can mix both patterns freely: several independent pools, each splitting its own share of the stream.
When you create a subscription, pick its starting point. earliest reads every record still retained. latest starts after the newest record that exists when you create it. Consumers that share a subscription get no guarantee about processing order between workers.
Standing up a stream, a producer and a consumer
You need a Cloudflare account on the Workers Paid plan, your account ID, curl and jq. K2 isn’t available on the Workers Free plan. The announcement says you can also create streams with cf, Wrangler or the dashboard. The steps below use the REST API, because it is what the docs show in full.
1. Create an API token with the K2 permissions
In the dashboard, open Account API tokens, choose Create Token, then Create Custom Token. Add K2 Config Write, K2 Produce and K2 Consume for your account. In production, give each service only what it needs: producers get K2 Produce, consumers get K2 Consume, and K2 Config Read is enough to list streams.
export ACCOUNT_ID=<YOUR_ACCOUNT_ID>
export CLOUDFLARE_API_TOKEN=<YOUR_API_TOKEN>
2. Create the stream with authentication on
curl "https://api.cloudflare.com/client/v4/accounts/$ACCOUNT_ID/k2/streams" \
--request POST \
--header "Authorization: Bearer $CLOUDFLARE_API_TOKEN" \
--header "Content-Type: application/json" \
--data '{
"name": "transactions",
"retention_seconds": 604800,
"http": { "enabled": true, "authentication": true }
}'
The response holds the stream id and its endpoint. The Workers binding input is on by default. Save both values:
export STREAM_ID=<STREAM_ID>
export K2_ENDPOINT=https://$STREAM_ID.k2.cloudflarestorage.com
Stream names allow letters, numbers and underscores, up to 128 characters, and you can’t rename a stream later.
3. Produce a batch over HTTP
Record content must be standard base64 with correct padding:
printf '%s' '{"order_id":1001,"status":"created"}' | base64
curl "$K2_ENDPOINT/produce" \
--request POST \
--header "Authorization: Bearer $CLOUDFLARE_API_TOKEN" \
--header "Content-Type: application/json" \
--data '{
"records": [
{
"content": "eyJvcmRlcl9pZCI6MTAwMSwic3RhdHVzIjoiY3JlYXRlZCJ9",
"headers": { "event-type": "order.created", "event-id": "ord-1001-created" }
}
]
}'
{ "success": true }
A batch is atomic: either every record is stored or none are. To save bandwidth, gzip the body and set Content-Encoding: gzip. The size limit applies both before and after decompression.
4. Produce from a Worker with a binding
A binding needs no API token. Add it to your Wrangler config, then run npx wrangler types:
{
"name": "orders-producer",
"main": "src/index.ts",
"compatibility_date": "2026-10-01",
"k2": [{ "binding": "ORDERS", "stream": "<STREAM_ID>" }]
}
export default {
async fetch(request, env): Promise<Response> {
const event = { order_id: 1001, status: "created" };
const result = await env.ORDERS.send([
{
content: new TextEncoder().encode(JSON.stringify(event)),
headers: { "event-type": "order.created", "event-id": crypto.randomUUID() },
},
]);
if (!result.success) {
console.error(`Produce failed: ${result.error.message}`);
return new Response("Failed to record event", {
status: result.error.retryable ? 503 : 500,
});
}
return new Response("Event recorded");
},
} satisfies ExportedHandler<Env>;
The stream must be in the same account as the Worker. On deploy, Cloudflare checks the binding and fails the deploy with code 10399 (bad or missing stream ID), 10400 (stream not in this account), 10401 (no permission) or 10402 (validation failed, try again).
5. Create one subscription per consumer group
curl "$K2_ENDPOINT/subscriptions" \
--request POST \
--header "Authorization: Bearer $CLOUDFLARE_API_TOKEN" \
--header "Content-Type: application/json" \
--data '{ "name": "analytics", "start_at": { "type": "earliest" } }'
Repeat with fraud-detection for the second group. If you send the same name with the same settings, K2 returns the existing ID, so this call is safe to keep in a deploy script.
6. Run a consumer loop that acks, nacks and backs off
This Node script polls the subscription, decodes records, skips duplicate event IDs, parks records that fail, acks the batch, backs off when the stream is empty, and logs lag on every batch. Run one copy per worker, each with a unique WORKER_ID.
// consumer.mjs
const endpoint = process.env.K2_ENDPOINT;
const token = process.env.CLOUDFLARE_API_TOKEN;
const sub = process.env.SUBSCRIPTION_ID;
const workerId = process.env.WORKER_ID; // unique per running process
const RETRYABLE = [10211, 10214, 10216, 10217];
async function call(path, body) {
const res = await fetch(`${endpoint}/subscriptions/${sub}${path}`, {
method: "POST",
headers: { Authorization: `Bearer ${token}`, "Content-Type": "application/json" },
body: JSON.stringify(body),
});
return res.json();
}
async function handle(event, headers) {
// Your processing logic.
}
// Replace these with a durable store (KV, D1, Redis) in production.
const seen = new Set();
const alreadyProcessed = async (id) => seen.has(id);
const markProcessed = async (id) => { seen.add(id); };
// Park records that fail so one bad record can't block the batch.
const park = async (failed) => console.error(JSON.stringify({ parked: failed }));
const sleep = (ms) => new Promise((r) => setTimeout(r, ms));
let idle = 1000;
while (true) {
const res = await call("/consume", { worker_id: workerId, max_records: 500 });
if (!res.success) {
const code = res.errors[0]?.code;
if (RETRYABLE.includes(code)) {
await sleep(idle);
idle = Math.min(idle * 2, 30000);
continue;
}
throw new Error(`consume failed: ${code} ${res.errors[0]?.message}`);
}
const { batch_id, records } = res.result;
if (!batch_id) {
await sleep(idle);
idle = Math.min(idle * 2, 30000);
continue;
}
idle = 1000;
const lagMs = Date.now() - records[records.length - 1].timestamp_ms;
console.log(JSON.stringify({ batch_id, records: records.length, lag_ms: lagMs }));
const failed = [];
for (const r of records) {
const id = r.headers?.["event-id"];
if (id && (await alreadyProcessed(id))) continue; // at-least-once: skip duplicates
try {
const event = JSON.parse(Buffer.from(r.content, "base64").toString("utf8"));
await handle(event, r.headers ?? {});
if (id) await markProcessed(id);
} catch (err) {
failed.push({ id, headers: r.headers, error: String(err) });
}
}
if (failed.length) await park(failed);
await call(`/batches/${batch_id}/ack`, { worker_id: workerId });
}
The loop acks every batch and parks records that fail, so one bad record never blocks the subscription. Save nack for failures that hit the whole batch, such as a downstream outage, so K2 redelivers it later.
max_records is a ceiling, from 1 to 10,000, and a short batch doesn’t mean the stream is drained. An empty result has batch_id set to null, which is your cue to wait before polling again.
7. Extend the lease for slow batches
If a batch may take longer than 5 minutes, extend the lease before it runs out. Each extension sets expiry to 5 minutes from the request:
curl "$K2_ENDPOINT/subscriptions/$SUBSCRIPTION_ID/batches/$BATCH_ID/extend" \
--request POST \
--header "Authorization: Bearer $CLOUDFLARE_API_TOKEN" \
--header "Content-Type: application/json" \
--data '{ "worker_id": "worker-1" }'
If the lease has already expired, been released or moved to another worker, K2 returns 409 with code 10218. Retrying can’t succeed. Stop work on that batch and request a new one: K2 sends the records again under a new batch_id.
The delivery rules your consumers have to handle
These come straight from the K2 docs. Each one changes how you write a consumer.
- At-least-once delivery. A record comes back if a lease expires, if you nack, or if a producer retries. K2 does not deduplicate. Give each record a stable ID in its headers, as in the examples above, and skip IDs you have already processed.
- Some produce failures have an unknown outcome. If
retryableis true (code 10211), the batch was not stored, so retry it with exponential backoff. Code 10212 means the batch may or may not have been stored. A retry can store it twice, which is one more reason for consumer-side dedupe. - No ordering across workers in one subscription. Records in the log are ordered, but workers that share a subscription can process them in any order. If order matters per entity, wait for message keys or design the handler to tolerate reordering.
- Acks and nacks cover whole batches. One bad record nacks the whole batch, and the subscription won’t move past those records until a later batch with them is acked. The docs describe no dead-letter queue or retry limit, so one poison record can keep coming back. Catch per-record failures, park the bad record somewhere you control (a second stream, R2 or a log line), and ack the rest.
- The server sets timestamps.
timestamp_msis when K2 received the batch. You can’t override it, so carry your own event time inside the content or a header. - Retention isn’t an exact delete. K2 removes expired records in the background, and they can stay readable for a while. Don’t use retention to meet a hard deletion deadline.
Beta limits and the pricing Cloudflare expects
During the public beta, each account can store up to 10 GB across all streams, and each stream accepts up to 30 MB/s of produce. To go higher, ask on the Cloudflare Discord or use the limit increase form linked from the limits page. Higher retention is also available on request.
The main K2 limits from the Cloudflare docs and announcement
| Limit | Value |
|---|---|
| Storage per account (beta) | 10 GB |
| Produce throughput per stream (beta) | 30 MB/s |
| Streams per account | 20 |
| Retention | 1 hour to 30 days, default 7 days |
| Produce request size | 5 MB |
| Record size, content plus headers | about 1 MB |
| Headers per record | 32 |
| Header name, value, total | 256 bytes, 8 KiB, 64 KiB |
| max_records per consume request | 10,000 |
| Data per consume response | 10 MB |
| Active leases per subscription | 128 |
| Subscriptions per stream | 100 |
| Lease length | 5 minutes |
| CORS origins per stream | 5 |
K2 usage isn’t billed during the beta. When billing starts, Cloudflare expects to charge $0.04 per GB produced, $0.04 per GB consumed and $0.02 per GB per month retained. The announcement doesn’t say how consumption is counted across subscriptions. A likely reading is that two subscriptions that each read every record count as two lots of data consumed, so budget for fan-out on that basis until Cloudflare confirms.
What goes wrong when teams first adopt K2
- The produce endpoint is public by accident. If you leave out
authenticationor set it to false, anyone who knows the stream ID can write to it. Set"authentication": truefor server-side producers. Browser producers need authentication off, because a token in page code is readable by anyone. In that case, list your origins incors.originsand validate every record in the consumer. - URL-safe base64. The HTTP API only accepts standard base64 with padding and rejects anything else with code 10204. Encode with
base64ortoBase64(). Don’t use URL-safe helpers. - Strings sent through the binding.
send()only takesArrayBufferorUint8Arraycontent, the same type for every record in a batch. Encode text withTextEncoderfirst, base64 strings included. - Assuming
send()throws. It returns a result object when K2 rejects a batch. Checkresult.successafter every call, or failures disappear silently. - Retrying everything. Only retry produce errors where
retryableis true. On the consume side, only retry 10211, 10214, 10216 and 10217, with backoff. Anything else needs a change to the request first. - Two processes with one
worker_id. Eachworker_idholds one lease at a time, and a second request from it gets the same batch back. Two processes sharing an ID end up doing the same work. Generate the ID from the host or pod name. - Running out of leases. A subscription allows 128 active leases. Past that, consume returns 429 with code 10216. That is the ceiling on parallel workers in one group.
- Changing a subscription. Subscriptions are immutable. Reusing a name with different settings returns 422 (code 10201). Create a new subscription and delete the old one, keeping in mind the cap of 100 per stream (code 10219).
- PATCH wipes input settings. Updating
httpreplaces the whole object. To add a CORS origin, sendenabledandauthenticationagain, or authentication goes back to its default. - Expecting low latency. About 1 second at p99 on produce is the price of building on R2. A request path that waits on a K2 write gets that second added. An Express tier with lower latency is on the roadmap.
Watching a K2 pipeline once it carries production traffic
I built and led the Workers observability team at Cloudflare, and for any buffer like this the first number I would chart is consumer lag. K2 keeps data safe while consumers are down, but only until retention expires. If lag gets close to retention_seconds, you lose records, and the K2 docs describe no warning for this. The lag_ms log line in the consumer loop above gives you that number per batch, per subscription.
The other signals worth an alert:
- Produce failures by code. A rise in 10211 means K2 is having trouble and your retries are holding. Any 10212 means possible duplicates. A 10206 or 10207 means a producer is sending requests over 5 MB or records over 1 MB.
- Nack rate per subscription. A steady nack rate usually means one poison record looping.
- 10218 on extend. Batches are taking longer than their leases, so work is being done twice.
- 10216 on consume. The group has hit 128 leases. More workers won’t help.
- Storage against the 10 GB beta cap. Long retention on a busy stream fills it fast.
What is on the K2 roadmap
Cloudflare lists these for the coming months: write parallelism up to multi-GB/s streams, message keys with key-based ordering, push-based Worker consumers, an Express tier with lower produce and end-to-end latency, and drop-in support for Apache Kafka clients. Cloudflare says a technical deep dive on K2’s design is coming. Until these features ship, plan around polling consumers, batch-level acks and no per-key ordering.
Running on Cloudflare? See how Polylane monitors Cloudflare in production.
Common questions.
When was Cloudflare K2 announced, and can I use it now?
Cloudflare announced K2 on 1 October 2026 and opened it as a public beta the same day. It is available to accounts on the Workers Paid plan. It isn't available on the Workers Free plan.
How much does Cloudflare K2 cost?
K2 usage isn't billed during the beta. Once billing starts, Cloudflare expects to charge $0.04 per GB produced, $0.04 per GB consumed and $0.02 per GB per month retained. The announcement doesn't say how consumption is counted across subscriptions.
Does K2 guarantee exactly-once delivery?
No. K2 delivers at least once and doesn't deduplicate records. Duplicates come from expired leases, nacks and producer retries, including code 10212, where a batch may or may not have been stored. Put a stable event ID in each record's headers and skip IDs your consumer has already processed.
How is K2 different from Cloudflare Queues?
Queues track individual items of work and offer retries, delays and dead-letter queues for each message. K2 is an ordered log built for high-volume data movement, long retention and fan-out, with records produced and consumed in batches. Use Queues for jobs and K2 for event streams that several systems read.
Can I point my Kafka clients at K2?
Not yet. Drop-in support for Apache Kafka clients is on the K2 roadmap, along with message keys and key-based ordering. Today you produce over the HTTP produce endpoint or a Workers binding, and you consume by polling the subscription's consume endpoint.
How long does K2 keep records?
Retention is set per stream, from 1 hour to 30 days, with a default of 7 days. Higher retention is available on request. Expired records are deleted in the background and can stay readable for a while, so don't rely on retention for exact deletion.
How fast is K2 produce latency?
At launch, a produce call takes about 1 second at the 99th percentile. The delay comes from batching writes in memory and storing them as segment files on R2. Cloudflare lists an Express tier with lower produce and end-to-end latency on the roadmap.
How many consumers can read one K2 stream?
A stream can have up to 100 subscriptions, and each one receives every record. Within a subscription, up to 128 leases can be active at once, and each worker_id holds one lease. That puts the ceiling at 128 parallel workers per consumer group.
Sources
Boris Tane is the founder of Polylane. He previously founded Baselime, observability for the future of the cloud, which Cloudflare acquired. At Cloudflare he built and led the Workers observability team.
Related
- How to Debug Cloudflare Workers Errors: Logs, Traces and Error Codes
Debug Cloudflare Workers errors step by step: read 1101 and 1102 codes, enable Workers Logs and source maps, use wrangler tail, DevTools and local traces.
- Cloudflare Durable Objects pending I/O keep-alive guide
From 2026-10-01, pending I/O keeps Cloudflare Durable Objects in memory after the client leaves. What counts, the 15-minute limit, flags and billing.
- Turn your app into a context graph
Agents that run software need a context graph of the app: every resource, what it connects to, the repository that deploys it and the team that owns it. How we built one on Durable Objects, keep it fresh across every provider without polling AWS, decide what counts as a change, and delete from it safely.
- Cloudflare Containers agent sandboxes: startup and setup
Cloudflare Containers now start agent sandboxes in a median 648 ms. Set up the durable_object policy, runtime images and snapshots, and know the limits.
- Fix “Durable Object reset because its code was updated”
Why a deploy makes Cloudflare Durable Objects throw “reset because its code was updated”, and how to retry safely, keep clients connected and lose no state.