Skip to content

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")
src/Orders.ts
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 }); // authenticated
Cloudflare.K2.Stream("Orders", { http: { authentication: false } }); // public
Cloudflare.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.

In the Worker’s init phase, ask Cloudflare.K2.WriteStream for a client. Its send appends a batch of records atomically.

src/Api.ts
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.

Declare the stream with an Effect Schema, and every binding on it is typed by that schema:

src/Orders.ts
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.

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.

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.

The producer code stays the same on every host. Only the layer changes:

  • WriteStreamBinding and StreamSinkBinding run in Cloudflare Workers and send through the native k2 binding.
  • WriteStreamHttp and StreamSinkHttp run 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.
  • WriteStreamLocal and StreamSinkLocal run 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).

Cloudflare.K2.consumeStreamRecords(stream, handler) subscribes a host to a stream. The handler receives each batch as an Effect Stream:

src/Fraud.ts
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 }.

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.

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.

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 Sink model behind StreamSink.

Reference:

  • K2 API reference: Stream, Subscription, WriteStream, StreamSink, ReadSubscription, and the polling layer.