Skip to content

Pipelines

Cloudflare Pipelines ingests events into R2. Producers send JSON events to a stream. A pipeline filters or reshapes them with SQL and writes them to a sink: plain JSON or Parquet files in a bucket, or rows in an Apache Iceberg table.

Pipelines is part of Basin, Cloudflare’s analytics stack on R2, and its resources are also exported as Cloudflare.Basin.*. Reach for Pipelines when events should land as files or tables you query later. Reach for K2 or Queues when your own code consumes the events.

By the end of this page, a Worker sends page views to a stream, and a pipeline writes them to an Iceberg table.

Declare the stream with an Effect Schema. Alchemy converts it into the stream’s field list, which Cloudflare enforces.

src/Events.ts
import * as Cloudflare from "alchemy/Cloudflare";
import * as Schema from "effect/Schema";
export const PageView = Schema.Struct({
userId: Schema.String,
path: Schema.String,
at: Schema.Date,
tags: Schema.optional(Schema.Array(Schema.String)),
});
export const PageViews = Cloudflare.Pipelines.Stream("PageViews", {
schema: PageView,
});

Arrays become list fields, nested structs become struct fields, and Schema.optional fields are not required. A schema with no stream representation, such as a union, fails when the program is built.

A stream’s schema can’t change after it is created. Changing it replaces the stream, which drops events it hasn’t processed yet and deletes the pipelines that read from it. Writing the same schema as an Effect Schema or as an explicit { fields } list is the same stream and causes no change.

Cloudflare.Pipelines.WriteStream is typed by the stream’s schema. The client encodes each value with the schema before sending it.

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 { PageViews } from "./Events.ts";
export default class Api extends Cloudflare.Worker<Api>()(
"Api",
{ main: import.meta.url },
Effect.gen(function* () {
const views = yield* Cloudflare.Pipelines.WriteStream(PageViews);
return {
fetch: Effect.gen(function* () {
const request = yield* HttpServerRequest;
yield* views.send([{ userId: "anon", path: request.url, at: new Date() }]);
return HttpServerResponse.text("ok");
}),
};
}).pipe(Effect.provide(Cloudflare.Pipelines.WriteStreamBinding)),
) {}

Cloudflare accepts events that don’t match a stream’s schema and drops them later during processing, so the client-side encoding is the only validation.

To drain an Effect Stream into the pipeline, use Cloudflare.Pipelines.StreamSink. It packs events into requests of up to 5 MB:

const sink = yield* Cloudflare.Pipelines.StreamSink(PageViews);
yield* Stream.fromIterable(events).pipe(Stream.run(sink));

Provide Cloudflare.Pipelines.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 pipelines binding.
  • WriteStreamHttp and StreamSinkHttp run on any host. They send to the stream’s HTTP endpoint with a scoped “Pipelines Send” token that Alchemy mints at deploy.
  • WriteStreamLocal and StreamSinkLocal run in scripts and Actions. They send to the stream’s HTTP endpoint with your current Cloudflare credentials.

The HTTP and local layers need the stream’s HTTP endpoint turned on.

A stream gets an HTTP ingest endpoint only when you ask for one:

Cloudflare.Pipelines.Stream("PageViews", { http: true }); // authenticated
Cloudflare.Pipelines.Stream("PageViews", {
http: { enabled: true, authentication: false }, // public, e.g. for browsers
});

A catalog sink writes the pipeline’s output to an Iceberg table. It needs a bucket with the catalog enabled. Declare them in your stack’s Effect, and pass the catalog resource to the sink so it is created after the catalog is enabled:

alchemy.run.ts
import * as Alchemy from "alchemy";
import * as Cloudflare from "alchemy/Cloudflare";
import * as Config from "effect/Config";
import * as Effect from "effect/Effect";
import { PageViews } from "./src/Events.ts";
export default Alchemy.Stack(
"Analytics",
{
providers: Cloudflare.providers(),
state: Cloudflare.state(),
},
Effect.gen(function* () {
const lake = yield* Cloudflare.R2.Bucket("Lake");
const catalog = yield* Cloudflare.R2.DataCatalog("Catalog", {
bucket: lake,
});
const views = yield* PageViews;
const table = yield* Cloudflare.Pipelines.Sink("PageViewTable", {
type: "basin_catalog",
catalog,
table: { namespace: "web", name: "page_views" },
token: yield* Config.Redacted("CATALOG_TOKEN"),
});
}),
);

The sink creates the namespace and table on its first write. The token is a Cloudflare API token with access to R2 and the catalog. Rotating it updates the sink without replacing it.

Add the pipeline in the same Effect:

alchemy.run.ts
import * as Output from "alchemy/Output";
const table = yield* Cloudflare.Pipelines.Sink("PageViewTable", { /* ... */ });
yield* Cloudflare.Pipelines.Pipeline("PageViewEtl", {
sql: Output.interpolate`
INSERT INTO ${table.name}
SELECT userId, path, at FROM ${views.name}
WHERE path NOT LIKE '/health%'`,
});

Alchemy validates the SQL with Cloudflare during plan, so a bad statement fails before anything is replaced. Reformatting the SQL is not a change.

An r2 sink writes events to a bucket as JSON or Parquet files, without a catalog:

const archive = yield* Cloudflare.Pipelines.Sink("EventArchive", {
type: "r2",
config: {
bucket: lake.bucketName,
credentials: {
accessKeyId: yield* Config.Redacted("R2_ACCESS_KEY_ID"),
secretAccessKey: yield* Config.Redacted("R2_SECRET_ACCESS_KEY"),
},
path: "events",
rollingPolicy: { intervalSeconds: 60 },
},
format: { type: "parquet", compression: "zstd" },
});

rollingPolicy decides when the sink closes a file and starts the next one: after intervalSeconds, after inactivitySeconds without new events, or when the file reaches fileSizeBytes. Point a pipeline at it the same way, with ${archive.name} in the SQL. Unlike a catalog sink, an r2 sink can be replaced freely.

Pipelines has no event source of its own, because its consumer is the sink. To process the files a sink writes, listen for object events on the bucket with Cloudflare.R2.BucketEventNotification.

Related:

  • Iceberg tables: manage the tables a catalog sink writes, and query them.
  • R2: the bucket a sink writes to.
  • K2 streams: a durable log your own code consumes.

Reference: