K2 streams
A K2 stream is a durable, ordered log of records stored on R2. Producers append records; readers consume them through subscriptions. Each subscription tracks its own position, so many readers can read the same records independently, and records stay in the stream for its retention period whether or not anyone has read them.
Reach for K2 when several consumers need the same events, when you need to replay history, or when you want a log you can read at your own pace. Reach for Queues when each message is a unit of work that one consumer handles once, with retries and a dead-letter queue.
| Queues | K2 | |
|---|---|---|
| Model | Work queue: a message is removed once it is acknowledged | Log: records stay for the retention period |
| Readers | One consumer | Any number of subscriptions, each reading every record |
| Delivery | Pushed to a consumer Worker | Pulled in leased batches over HTTP |
| Consume from a Worker | Yes | Not yet |
| Replay | No | From the start of retention (startAt: "earliest") |
Create a stream
Section titled “Create a stream”import * as Cloudflare from "alchemy/Cloudflare";
export const Orders = Cloudflare.K2.Stream("Orders", { retention: "7 days",});The stream name is generated from the stack, stage, and logical ID. Changing
retention updates the stream in place. It accepts any Duration.Input
between one hour and 30 days, and defaults to seven days.
A stream accepts records from Workers by default. Its HTTP input is off unless you turn it on:
Cloudflare.K2.Stream("Orders", { http: true }); // authenticatedCloudflare.K2.Stream("Orders", { http: { authentication: false } }); // publicCloudflare.K2.Stream("Orders", { http: { cors: ["https://shop.dev"] } });Turn on the HTTP input when a host outside Cloudflare produces to the stream,
or when you produce with WriteStreamHttp.
Produce from a Worker
Section titled “Produce from a Worker”In the Worker’s init phase, ask Cloudflare.K2.WriteStream for a client. Its
send appends a batch of records atomically.
import * as Cloudflare from "alchemy/Cloudflare";import * as Effect from "effect/Effect";import { HttpServerRequest } from "effect/http/HttpServerRequest";import * as HttpServerResponse from "effect/http/HttpServerResponse";import { Orders } from "./Orders.ts";
export default class Api extends Cloudflare.Worker<Api>()( "Api", { main: import.meta.url }, Effect.gen(function* () { const orders = yield* Cloudflare.K2.WriteStream(Orders);
return { fetch: Effect.gen(function* () { const request = yield* HttpServerRequest; const body = yield* request.text; yield* orders.send([{ content: body, headers: { source: "api" } }]); return HttpServerResponse.empty({ status: 202 }); }), }; }).pipe(Effect.provide(Cloudflare.K2.WriteStreamBinding)),) {}A record’s content is a string (sent as UTF-8), a Uint8Array, or an
ArrayBuffer. headers is an optional map of up to 32 string pairs. A batch
can be up to 5 MB, and each record up to 1 MB.
WriteStreamBinding sends through the Worker’s native k2 binding, so it
needs no API token.
Type records with a schema
Section titled “Type records with a schema”Declare the stream with an Effect Schema, and every binding on it is typed by that schema:
import * as Cloudflare from "alchemy/Cloudflare";import * as Schema from "effect/Schema";
export const Order = Schema.Struct({ orderId: Schema.String, amountCents: Schema.Number,});
export const Orders = Cloudflare.K2.Stream("Orders", { retention: "7 days", schema: Order,});Now orders.send takes Order values. The client encodes each one as JSON
and sets a content-type: application/json header:
yield* orders.send([{ orderId: "o_1", amountCents: 4200 }]);K2 itself stores bytes, so the schema is a client-side codec. Changing it never updates or replaces the stream. Producers and consumers read the schema from the same stream, so they can’t disagree about the record type.
Handle produce errors
Section titled “Handle produce errors”Failures are typed. Two of them tell you whether the batch was stored:
yield* orders.send(batch).pipe( // The batch was not stored. Retrying it is safe. Effect.catchTag("K2Unavailable", () => retryLater(batch)), // The batch may or may not have been stored, and K2 does not deduplicate. // Alchemy never retries this one for you. Effect.catchTag("K2AppendOutcomeUnknown", () => reconcile(batch)),);The native binding and the HTTP path report the same error tags, so this code works with either layer.
Drain a Stream into K2
Section titled “Drain a Stream into K2”Cloudflare.K2.StreamSink exposes the stream as an Effect
Sink. It packs records into requests of
up to 5 MB and retries only the errors that mean the batch was not stored.
const sink = yield* Cloudflare.K2.StreamSink(Orders);yield* Stream.fromIterable(pending).pipe(Stream.run(sink));Provide Cloudflare.K2.StreamSinkBinding alongside WriteStreamBinding.
Produce from outside a Worker
Section titled “Produce from outside a Worker”The producer code stays the same on every host. Only the layer changes:
WriteStreamBindingandStreamSinkBindingrun in Cloudflare Workers and send through the nativek2binding.WriteStreamHttpandStreamSinkHttprun on any host, such as Lambda or a container. They send through the stream’s HTTP input with a scoped “K2 Produce” token that Alchemy mints at deploy.WriteStreamLocalandStreamSinkLocalrun in scripts and Actions. They send through the stream’s HTTP input with your current Cloudflare credentials.
The HTTP and local layers need the stream’s HTTP input turned on
(http: true).
Consume from a long-running host
Section titled “Consume from a long-running host”Cloudflare.K2.consumeStreamRecords(stream, handler) subscribes a host to a
stream. The handler receives each batch as an Effect Stream:
import * as Cloudflare from "alchemy/Cloudflare";import * as Fly from "alchemy/Fly";import * as Effect from "effect/Effect";import * as Layer from "effect/Layer";import * as Stream from "effect/Stream";import { Orders } from "./Orders.ts";
export default class Fraud extends Fly.Service<Fraud>()( "Fraud", { main: import.meta.url, region: "iad" }, Effect.gen(function* () { yield* Cloudflare.K2.consumeStreamRecords( Orders, { startAt: "earliest", concurrency: 4 }, (records) => records.pipe( Stream.map((record) => record.value), Stream.filter((order) => order.amountCents > 100_000), Stream.runForEach((order) => Effect.log(`review ${order.orderId}`)), ), );
return {}; }).pipe( Effect.provide( Layer.provideMerge( Cloudflare.K2.StreamEventSourcePolling, Cloudflare.K2.ReadSubscriptionHttp, ), ), ),) {}At deploy time, consumeStreamRecords creates a subscription owned by the
host. At runtime, StreamEventSourcePolling runs concurrency pollers. Each
poller leases a batch of up to maxRecords (default 100) and runs the
handler:
- When the handler succeeds, the batch is acknowledged.
- When it fails, the failure is logged and the batch is released for redelivery.
- While the handler runs, the five-minute lease is extended every two minutes.
- When a batch comes back empty, the poller waits
idleBackoff(default one second) before it polls again.
ReadSubscriptionHttp reads with a scoped “K2 Consume” token that Alchemy
mints at deploy. StreamEventSourcePolling needs a host with a long-lived
process (Fly, ECS, EC2, Railway, Hetzner, Docker), so it doesn’t work in a
Worker.
For a stream with a schema, each record’s JSON content is decoded into
record.value. For a stream without one, each record is
{ content: Uint8Array, headers, timestamp }.
Share a subscription between hosts
Section titled “Share a subscription between hosts”By default every consuming host gets its own subscription, so each host sees
every record. To split records between hosts instead, declare a
Subscription and pass it to each one:
export const FraudWorkers = Cloudflare.K2.Subscription("FraudWorkers", { stream: Orders, startAt: "latest",});yield* Cloudflare.K2.consumeStreamRecords( Orders, { startAt: "earliest", concurrency: 4 }, { subscription: FraudWorkers, concurrency: 4 },Hosts sharing a subscription compete for batches. A subscription allows up
to 128 leases at a time. Subscriptions can’t be changed after they are
created, so changing startAt replaces the subscription.
Control leases yourself
Section titled “Control leases yourself”Cloudflare.K2.ReadSubscription is the client consumeStreamRecords uses
internally. Use it when you need your own loop:
const inbox = yield* Cloudflare.K2.ReadSubscription(FraudWorkers);
const batch = yield* inbox.consume({ workerId: "worker-1", maxRecords: 100 });if (Option.isSome(batch)) { yield* handle(batch.value.records).pipe( Effect.andThen(inbox.ack(batch.value)), Effect.catchCause(() => inbox.nack(batch.value)), );}consume returns Option.none() when there are no records. A worker that
already holds a lease gets the same batch back. extend(batch) renews a
lease for another five minutes, and fails with K2LeaseLost once the lease
is gone.
Where next
Section titled “Where next”Related:
- Queues: a work queue with push delivery to Workers, retries, and dead-letter queues.
- Pipelines: land events as files or Iceberg tables instead of consuming them yourself.
- Workers: the runtime that produces records.
- Sinks: the Effect
Sinkmodel behindStreamSink.
Reference:
- K2 API reference:
Stream,Subscription,WriteStream,StreamSink,ReadSubscription, and the polling layer.