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.
Create a typed stream
Section titled “Create a typed stream”Declare the stream with an Effect Schema. Alchemy converts it into the stream’s field list, which Cloudflare enforces.
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.
Send events from a Worker
Section titled “Send events from a Worker”Cloudflare.Pipelines.WriteStream is typed by the stream’s schema. The client
encodes each value with the schema before sending it.
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.
Send events from outside a Worker
Section titled “Send events 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 nativepipelinesbinding.WriteStreamHttpandStreamSinkHttprun on any host. They send to the stream’s HTTP endpoint with a scoped “Pipelines Send” token that Alchemy mints at deploy.WriteStreamLocalandStreamSinkLocalrun 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.
Turn on the HTTP endpoint
Section titled “Turn on the HTTP endpoint”A stream gets an HTTP ingest endpoint only when you ask for one:
Cloudflare.Pipelines.Stream("PageViews", { http: true }); // authenticatedCloudflare.Pipelines.Stream("PageViews", { http: { enabled: true, authentication: false }, // public, e.g. for browsers});Write events to an Iceberg table
Section titled “Write events to an Iceberg table”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:
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.
Connect the stream and sink with SQL
Section titled “Connect the stream and sink with SQL”Add the pipeline in the same Effect:
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.
Write files instead of a table
Section titled “Write files instead of a table”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.
React to new files
Section titled “React to new files”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.
Where next
Section titled “Where next”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:
- Pipelines API reference:
Stream,Sink,Pipeline,WriteStream, andStreamSink.