This site is not affiliated with or endorsed by Cloudflare, Inc. It simply showcases experiments built using Cloudflare services.
Cloudflare Experiments
Storage & Data

Event Pipeline

Ingest events into Cloudflare Pipelines (stream to R2/Iceberg)

Ingest events into Cloudflare Pipelines for streaming to R2/Iceberg. When the Pipelines binding is unavailable, events append as JSON Lines to an R2 bucket.

Features

  • POST /events - one event, an array (max 100), or { "events": [...] }
  • GET /events/sample - example schema and accepted body shapes
  • PIPELINE binding for Cloudflare Pipelines; EVENTS R2 bucket for fallback

API Reference

POST /events

Accepts any of:

  • A single event object
  • An array of event objects (max 100)
  • { "events": [ ... ] }

Example Request

curl -X POST "https://your-worker.workers.dev/events" \
  -H "Content-Type: application/json" \
  -d '{"events":[{"type":"page_view","timestamp":"2026-01-15T12:00:00.000Z","userId":"user_123","properties":{"path":"/home"}}]}'

Success Response

{
  "ok": true,
  "count": 1,
  "transport": "pipeline"
}

transport is "pipeline" when PIPELINE is bound, otherwise "r2" (appends to events/YYYY-MM-DD.jsonl).

GET /events/sample

Returns an example event schema and accepted request shapes.

Example Request

curl "https://your-worker.workers.dev/events/sample"

Error Codes

  • 400 - Invalid JSON, non-object events, or empty batch (INVALID_BODY)
  • 400 - More than 100 events (TOO_MANY_EVENTS)
  • 502 - Pipeline or R2 write failed (PIPELINE_ERROR)

Use Cases

  • Learn Cloudflare Pipelines ingestion from Workers
  • Stream analytics or product events toward R2/Iceberg
  • Fall back to R2 JSONL when Pipelines is unavailable locally
  • Prototype high-volume event APIs before adding auth

Limitations

  • Max 100 events per request
  • Requires a dashboard Pipeline and/or R2 bucket matching wrangler names
  • No schema validation beyond “JSON object(s)”
  • R2 fallback is append-only JSONL, not a query API

Use in your project

Copy these files into an existing Worker. Prefer Deployment to try the full experiment first. Source: apps/experiments/event-pipeline.

DependenciesNoneBindingsR2PlatformWorkers runtime
import { MAX_EVENTS } from "../constants/defaults";import type { EventObject } from "../types/events";import type { PipelineBinding } from "../types/env";export type ParseEventsResult =  | { ok: true; events: EventObject[] }  | { ok: false; code: "INVALID_BODY" | "TOO_MANY_EVENTS"; message: string };function isPlainObject(value: unknown): value is EventObject {  return typeof value === "object" && value !== null && !Array.isArray(value);}function isJsonSerializable(value: unknown): boolean {  try {    JSON.stringify(value);    return true;  } catch {    return false;  }}/** * Accepts a single object, an array of objects, or `{ events: object[] }`. */export function parseEventsBody(body: unknown): ParseEventsResult {  if (body === null || body === undefined) {    return {      ok: false,      code: "INVALID_BODY",      message: "Body must be a JSON object, array of objects, or { events: object[] }",    };  }  let raw: unknown[];  if (Array.isArray(body)) {    raw = body;  } else if (isPlainObject(body)) {    if (Array.isArray(body.events)) {      raw = body.events;    } else {      raw = [body];    }  } else {    return {      ok: false,      code: "INVALID_BODY",      message: "Body must be a JSON object, array of objects, or { events: object[] }",    };  }  if (raw.length > MAX_EVENTS) {    return {      ok: false,      code: "TOO_MANY_EVENTS",      message: `Too many events: max ${MAX_EVENTS} per request`,    };  }  if (raw.length === 0) {    return {      ok: false,      code: "INVALID_BODY",      message: "At least one event object is required",    };  }  const events: EventObject[] = [];  for (const item of raw) {    if (!isPlainObject(item)) {      return {        ok: false,        code: "INVALID_BODY",        message: "Each event must be a JSON object",      };    }    if (!isJsonSerializable(item)) {      return {        ok: false,        code: "INVALID_BODY",        message: "Each event must be JSON-serializable",      };    }    events.push(item);  }  return { ok: true, events };}function eventsJsonlKey(date = new Date()): string {  const yyyy = date.getUTCFullYear();  const mm = String(date.getUTCMonth() + 1).padStart(2, "0");  const dd = String(date.getUTCDate()).padStart(2, "0");  return `events/${yyyy}-${mm}-${dd}.jsonl`;}/** * Send events via Pipelines when bound; otherwise append JSON lines to R2. */export async function ingestEvents(  events: EventObject[],  pipeline: PipelineBinding | undefined,  bucket: R2Bucket): Promise<"pipeline" | "r2"> {  if (pipeline) {    await pipeline.send(events);    return "pipeline";  }  const key = eventsJsonlKey();  const existing = await bucket.get(key);  const previous = existing ? await existing.text() : "";  const lines = events.map((e) => JSON.stringify(e)).join("\n");  const next = previous ? `${previous}\n${lines}\n` : `${lines}\n`;  await bucket.put(key, next, {    httpMetadata: { contentType: "application/x-ndjson" },  });  return "r2";}

Deployment

Click the deploy button

Deploy to Cloudflare Workers

Configure Pipelines and R2

  1. Create a Pipeline named events-pipeline (or update wrangler.json)
  2. Create an R2 bucket event-pipeline-events for the fallback path
  3. Confirm bindings: PIPELINE → pipeline name, EVENTS → bucket

Test your deployment

curl -X POST "https://your-worker.workers.dev/events" \
  -H "Content-Type: application/json" \
  -d '{"type":"page_view","timestamp":"2026-01-15T12:00:00.000Z"}'

Local Development

cd apps/experiments/event-pipeline
npm install
npm run dev
curl -X POST "http://localhost:8787/events" \
  -H "Content-Type: application/json" \
  -d '{"type":"page_view","timestamp":"2026-01-15T12:00:00.000Z"}'

Without a Pipelines binding, wrangler dev uses the R2 fallback when EVENTS is available.

Configuration

wrangler.json declares:

  • Pipelines binding PIPELINE → events-pipeline
  • R2 bucket EVENTS → event-pipeline-events

Cloudflare Features Used

  • Workers - Edge compute runtime
  • Pipelines - Stream events to R2/Iceberg
  • R2 - JSONL fallback storage

On this page