Skip to content

Cloudflare.Pipelines reference

Source: src/Cloudflare/Pipelines/LegacyPipeline.ts

A legacy Cloudflare Pipeline — the original HTTP-ingest → R2 batch product (/accounts/{account}/pipelines).

A legacy pipeline accepts JSON events over HTTP (and/or a Worker pipelines binding) and batches them into an R2 bucket using S3-compatible credentials.

LegacyPipeline: Creating a Legacy Pipeline

Section titled “LegacyPipeline: Creating a Legacy Pipeline”

HTTP ingest into R2

The S3-compatible credentials are derived from a Cloudflare API token: the access key id is the token id and the secret is the SHA-256 hex digest of the token value.

const bucket = yield* Cloudflare.R2.Bucket("events", {});
const pipeline = yield* Cloudflare.Pipelines.LegacyPipeline("ingest", {
destination: {
bucket: bucket.bucketName,
credentials: {
accessKeyId: yield* Config.Redacted("R2_ACCESS_KEY_ID"),
secretAccessKey: yield* Config.Redacted("R2_SECRET_ACCESS_KEY"),
},
},
});
// POST events to pipeline.endpoint

Tuned batching and CORS

const pipeline = yield* Cloudflare.Pipelines.LegacyPipeline("ingest", {
source: [
{ type: "http", cors: { origins: ["https://example.com"] } },
],
destination: {
bucket: bucket.bucketName,
credentials,
batch: { maxDurationS: 10, maxRows: 1000 },
compression: "gzip",
prefix: "ingest",
},
});

Source: src/Cloudflare/Pipelines/Pipeline.ts

A Cloudflare SQL Pipeline — the transform of the Pipelines product. A pipeline is a single SQL statement that reads events from a Stream and writes them to a Sink, both referenced by name.

Also exported as Cloudflare.Basin.Pipeline.

The SQL is fixed at creation: changing it (or the name) triggers a replacement. Nothing references a pipeline downstream, so replacements are cheap. A changed statement is checked with Cloudflare’s SQL validator during plan, so a syntax error fails the plan (with PipelineSqlInvalid) instead of deleting the running pipeline. References to streams/sinks that do not exist yet are not an error at plan time — they are created in the same deploy.

Pipelines is ingest-only: to process data after the sink writes it, subscribe to the sink’s bucket with Cloudflare.R2.BucketEventNotification.

Stream → Sink passthrough

const stream = yield* Cloudflare.Basin.Stream("events", {});
const sink = yield* Cloudflare.Basin.Sink("events-sink", {
type: "r2",
config: { bucket: bucket.bucketName, credentials },
});
const pipeline = yield* Cloudflare.Basin.Pipeline("etl", {
sql: Output.interpolate`INSERT INTO ${sink.name} SELECT * FROM ${stream.name}`,
});

Filtering transform

const pipeline = yield* Cloudflare.Basin.Pipeline("errors-only", {
sql: Output.interpolate`INSERT INTO ${sink.name} SELECT * FROM ${stream.name} WHERE level = 'error'`,
});

Source: src/Cloudflare/Pipelines/Sink.ts

A Cloudflare Pipelines sink — the destination of the Pipelines product (also exported as Cloudflare.Basin.Sink). A SQL Pipeline reads events from a Stream and writes them to a sink, which stores them in R2 either as raw files (r2) or as an Iceberg table via the Basin Catalog (basin_catalog, a.k.a. r2_data_catalog).

Sinks have no update API: changing the destination (bucket, path, table, format, partitioning, rolling policy) triggers a replacement. Equivalent spellings never replace — basin_catalog ⇄ r2_data_catalog, the config form ⇄ the catalog / table form, omitted ⇄ explicit defaults. With engine-generated names a replacement is seamless (the new sink gets a fresh name before the old one is deleted); with an explicit name it collides, so prefer generated names.

Credentials (credentials, token) are write-only and are not part of the replacement diff, so rotating them never recreates the sink — a catalog sink could not be recreated anyway, because Cloudflare refuses to create sinks for existing Iceberg tables. The sink keeps the credentials it was created with; to move it onto new ones, rename it (which replaces it).

Pipelines is ingest-only — to react to files the sink writes, subscribe to its bucket with Cloudflare.R2.BucketEventNotification.

R2 sink with JSON output

The S3-compatible credentials are derived from a Cloudflare API token: the access key id is the token id and the secret is the SHA-256 hex digest of the token value.

const bucket = yield* Cloudflare.R2.Bucket("events", {});
const sink = yield* Cloudflare.Basin.Sink("events-sink", {
type: "r2",
config: {
bucket,
credentials: {
accessKeyId: yield* Config.Redacted("R2_ACCESS_KEY_ID"),
secretAccessKey: yield* Config.Redacted("R2_SECRET_ACCESS_KEY"),
},
path: "ingest",
rollingPolicy: { intervalSeconds: 30 },
},
});

Parquet output

const sink = yield* Cloudflare.Basin.Sink("parquet-sink", {
type: "r2",
config: { bucket, credentials },
format: { type: "parquet", compression: "zstd" },
});

Iceberg table sink

const bucket = yield* Cloudflare.R2.Bucket("Lakehouse", {});
const catalog = yield* Cloudflare.Basin.Catalog("Catalog", {
bucket,
});
const sink = yield* Cloudflare.Basin.Sink("PageViewTable", {
type: "basin_catalog",
catalog,
table: { namespace: "web", name: "page_views" },
token: yield* Config.Redacted("CATALOG_TOKEN"),
});

Bucket-name form

const sink = yield* Cloudflare.Basin.Sink("iceberg-sink", {
type: "r2_data_catalog",
config: {
bucket,
tableName: "events",
namespace: "default",
token: yield* Config.Redacted("CATALOG_TOKEN"),
},
});

Source: src/Cloudflare/Pipelines/Stream.ts

A Cloudflare Pipelines stream — the ingestion endpoint of the Pipelines product (also exported as Cloudflare.Basin.Stream). Events are sent to a stream from Workers (WriteStream / StreamSink) or over HTTP, transformed by a SQL Pipeline, and written to a Sink.

The stream’s schema and format are fixed at creation (changing them triggers a replacement, which drops buffered events and dependent pipelines); the HTTP endpoint and Worker-binding toggles are mutable in place.

Pipelines is ingest-only — there is no event source on a stream. To process data after it lands, subscribe to the sink’s bucket with Cloudflare.R2.BucketEventNotification.

Cloudflare accepts records that violate a structured stream’s schema (they are dropped later, during processing). Declare the schema as an Effect Schema to have the producer clients validate and encode every record before it is sent.

Unstructured stream with default settings

const stream = yield* Cloudflare.Basin.Stream("events", {});

Structured stream from an Effect Schema

class PageView extends Schema.Class<PageView>("PageView")({
url: Schema.String,
at: Schema.Date,
tags: Schema.Array(Schema.String),
user: Schema.optional(Schema.Struct({ id: Schema.String })),
}) {}
// Stream<PageView>: WriteStream / StreamSink are typed and encode records
const stream = yield* Cloudflare.Basin.Stream("PageViews", {
schema: PageView,
});

Structured stream with an explicit field list

const stream = yield* Cloudflare.Basin.Stream("clicks", {
schema: {
fields: [
{ type: "string", name: "url", required: true },
{ type: "timestamp", name: "ts", unit: "millisecond" },
{ type: "list", name: "tags", items: { type: "string" } },
{
type: "struct",
name: "user",
fields: [{ type: "string", name: "id" }],
},
],
},
});

Authenticated endpoint

const stream = yield* Cloudflare.Basin.Stream("events", { http: true });
// POST a JSON array to stream.endpoint with a `Pipelines Send` token,
// or use Cloudflare.Basin.WriteStreamHttp / WriteStreamLocal.

Public endpoint with CORS

const stream = yield* Cloudflare.Basin.Stream("beacons", {
http: {
enabled: true,
authentication: false,
cors: { origins: ["https://app.example.com"] },
},
});
const pipeline = yield* Cloudflare.Basin.Pipeline("etl", {
sql: Output.interpolate`INSERT INTO ${sink.name} SELECT * FROM ${stream.name}`,
});

Source: src/Cloudflare/Pipelines/StreamSink.ts

Binding service that exposes a Pipelines Stream as an Effect Sink: run a Stream of records into it and every element is ingested. Also exported as Cloudflare.Basin.StreamSink.

Records are typed, validated and encoded by the stream’s Effect Schema (plain JSON objects when it has none). Each upstream chunk is packed greedily into requests of at most 5 MB, preserving order.

Provide StreamSinkBinding (native Worker binding), StreamSinkHttp (scoped Pipelines Send token) or StreamSinkLocal (current credentials, for Actions) — each is the sink layered over the matching WriteStream implementation.

StreamSink: Draining a Stream into a Pipelines Stream

Section titled “StreamSink: Draining a Stream into a Pipelines Stream”

Run an Effect Stream into a Pipelines stream

export default Cloudflare.Worker(
"Api",
{ main: import.meta.url },
Effect.gen(function* () {
const sink = yield* Cloudflare.Basin.StreamSink(PageViews);
return {
fetch: Effect.gen(function* () {
yield* EffectStream.fromIterable(views).pipe(EffectStream.run(sink));
return HttpServerResponse.empty({ status: 202 });
}),
};
}).pipe(Effect.provide(Cloudflare.Basin.StreamSinkBinding)),
);

Controlling request size

// at most 500 records per send
yield* events.pipe(EffectStream.rechunk(500), EffectStream.run(sink));

Source: src/Cloudflare/Pipelines/StreamSinkBinding.ts Kind: Layer · Provides: Cloudflare.Pipelines.StreamSink

Implementation of the StreamSink service over the native Worker pipelines binding (WriteStreamBinding). Registers the same binding as WriteStream, so a Worker can use both on one stream.

Effect.gen(function* () {
const sink = yield* Cloudflare.Basin.StreamSink(PageViews);
// ...
}).pipe(Effect.provide(Cloudflare.Basin.StreamSinkBinding));

Source: src/Cloudflare/Pipelines/StreamSinkHttp.ts Kind: Layer · Provides: Cloudflare.Pipelines.StreamSink

Implementation of the StreamSink service over the stream’s HTTP ingest endpoint (WriteStreamHttp), authenticated with a scoped Pipelines Send API token bound into the host. The stream needs http enabled.

Effect.gen(function* () {
const sink = yield* Cloudflare.Basin.StreamSink(PageViews);
// ...
}).pipe(Effect.provide(Cloudflare.Basin.StreamSinkHttp));

Source: src/Cloudflare/Pipelines/StreamSinkLocal.ts Kind: Layer · Provides: Cloudflare.Pipelines.StreamSink

Implementation of the StreamSink service over the stream’s HTTP ingest endpoint with the current credentials (WriteStreamLocal) — for draining an Effect Stream into a Pipelines stream from an Action, with no Worker host. The stream needs http enabled.

const Seed = Alchemy.Action(
"Seed",
Effect.gen(function* () {
const sink = yield* Cloudflare.Basin.StreamSink(PageViews);
return Effect.fn(function* () {
yield* EffectStream.fromIterable(views).pipe(EffectStream.run(sink));
});
}).pipe(Effect.provide(Cloudflare.Basin.StreamSinkLocal)),
);

Source: src/Cloudflare/Pipelines/WriteStream.ts

Binding service that turns a Pipelines Stream (or a LegacyPipeline) into a WriteStreamClient you can call from a Worker’s (or Action’s) runtime Effect. Also exported as Cloudflare.Basin.WriteStream.

Producer-only — send ingests a batch of records into the stream. When the stream was declared with an Effect Schema, the client is typed by it (send(records: PageView[])) and every record is validated and encoded with Schema.toCodecJson before it is sent — Cloudflare accepts schema-violating records and silently drops them later, so this is the only validation. Without a schema, send takes plain JSON objects.

Provide one implementation layer:

  • WriteStreamBinding — native Worker pipelines binding
  • WriteStreamHttp — the stream’s HTTP ingest endpoint with a scoped Pipelines Send API token (the stream needs http enabled)
  • WriteStreamLocal — the HTTP ingest endpoint with the current credentials, for Actions and other deploy-time Effects

Typed producer route

class PageView extends Schema.Class<PageView>("PageView")({
url: Schema.String,
at: Schema.Date,
}) {}
export const PageViews = Cloudflare.Basin.Stream("PageViews", {
schema: PageView,
});
export default Cloudflare.Worker(
"Api",
{ main: import.meta.url },
Effect.gen(function* () {
const views = yield* Cloudflare.Basin.WriteStream(PageViews);
return {
fetch: Effect.gen(function* () {
yield* views.send([new PageView({ url: "/", at: new Date() })]);
return HttpServerResponse.empty({ status: 202 });
}),
};
}).pipe(Effect.provide(Cloudflare.Basin.WriteStreamBinding)),
);

Untyped stream

const events = yield* Cloudflare.Basin.WriteStream(Events);
yield* events.send([{ event: "click", at: new Date().toISOString() }]);
yield* views.send(batch).pipe(
Effect.catchTag("StreamSendError", (e) =>
Effect.logWarning(`dropped ${batch.length} events: ${e.message}`),
),
);

Source: src/Cloudflare/Pipelines/WriteStreamBinding.ts Kind: Layer · Provides: Cloudflare.Pipelines.WriteStream

Implementation of the WriteStream service over a native Worker pipelines binding. Registers the binding on the host Worker and, when the stream has an Effect Schema, encodes records with it before send.

Effect.gen(function* () {
const views = yield* Cloudflare.Basin.WriteStream(PageViews);
// ...
}).pipe(Effect.provide(Cloudflare.Basin.WriteStreamBinding));

Source: src/Cloudflare/Pipelines/WriteStreamHttp.ts Kind: Layer · Provides: Cloudflare.Pipelines.WriteStream

HTTP-backed implementation of the WriteStream service. Creates a scoped AccountApiToken with the Pipelines Send permission, binds it into the host, and POSTs records to the stream’s HTTP ingest endpoint (https://{stream_id}.ingest.cloudflare.com).

The stream must have its HTTP endpoint enabled — declare it with http: true (authenticated). Records are validated and encoded with the stream’s Effect Schema exactly as with WriteStreamBinding.

const Events = Cloudflare.Basin.Stream("Events", { http: true });
Effect.gen(function* () {
const events = yield* Cloudflare.Basin.WriteStream(Events);
// ...
}).pipe(Effect.provide(Cloudflare.Basin.WriteStreamHttp));

Source: src/Cloudflare/Pipelines/WriteStreamLocal.ts Kind: Layer · Provides: Cloudflare.Pipelines.WriteStream

Local implementation of the WriteStream service — POSTs records to the stream’s HTTP ingest endpoint with the current credentials instead of a native Worker binding (WriteStreamBinding) or a scoped API token (WriteStreamHttp).

Provide it on an Action (or any deploy-time Effect) to send events with the same typed client you use inside a Worker. The stream needs its HTTP endpoint enabled (http: true), and the current credentials need the Pipelines Send permission when it is authenticated.

const Seed = Alchemy.Action(
"Seed",
Effect.gen(function* () {
const views = yield* Cloudflare.Basin.WriteStream(PageViews);
return Effect.fn(function* () {
yield* views.send([new PageView({ url: "/", at: new Date() })]);
});
}).pipe(Effect.provide(Cloudflare.Basin.WriteStreamLocal)),
);