All skills
cloudflare avatar

/basin

@41e0d19 official
by cloudflarecloudflare/skills3k stars
298

Build and troubleshoot Cloudflare Basin analytics workflows with Basin Pipelines, Basin Catalog, and Basin SQL. Use for streaming data into R2 Iceberg tables, managing catalogs, or querying those tables; also use for requests using the former Data Platform, Pipelines, R2 Data Catalog, or R2 SQL names.

Use this Skill: https://skilld.dev/gh/cloudflare/skills/basin

This session only. Nothing lands on disk.

referencespipelinespatterns.md

≈1.2k tokens on demand. Your agent reads this file only when SKILL.md points to it.

Basin Pipelines Patterns

Code-first patterns. For observability dataset/field schemas and Logpush dataset lists, pull https://developers.cloudflare.com/basin-pipelines/observability/metrics/index.md and https://developers.cloudflare.com/basin-pipelines/streams/logpush/index.md.

Fire-and-Forget Producer

export default {
  async fetch(req, env, ctx) {
    const event = { event_id: crypto.randomUUID(), event_type: "page_view", timestamp: new Date().toISOString() };
    ctx.waitUntil(env.MY_STREAM.send([event]));  // don't block the response
    return new Response("OK");
  }
};

Client-Side Validation with Zod

Structured streams drop invalid events silently during processing. Validate before sending for immediate feedback.

import { z } from "zod";

const EventSchema = z.object({
  event_id: z.string(),
  category: z.enum(["purchase", "view"]),
  amount: z.number().positive().optional(),
});

const validated = EventSchema.parse(rawEvent);  // throws synchronously
await env.MY_STREAM.send([validated]);

Scheduled Collector Worker

// wrangler.jsonc
{
  "name": "collector",
  "pipelines": [{ "stream": "<STREAM_ID>", "binding": "EVENT_STREAM" }],
  "triggers": { "crons": ["*/5 * * * *"] }
}
export default {
  async scheduled(event, env, ctx) {
    const items = await (await fetch("https://api.example.com/data")).json();
    const events = items.map(i => ({
      event_id: crypto.randomUUID(),
      timestamp: new Date().toISOString(),
      category: i.type, amount: i.value,
    }));
    await env.EVENT_STREAM.send(events);
  },
};

Logpush → Basin Pipelines

Basin Pipelines is a native Logpush destination — ingest Cloudflare logs, transform with SQL, store as Iceberg/Parquet. For the current supported dataset list and field names, pull the Logpush doc above.

INSERT INTO http_logs_sink
SELECT
  ClientIP,
  EdgeResponseStatus,
  to_timestamp_micros(EdgeStartTimestamp) AS event_time,
  upper(ClientRequestMethod) AS method,
  sha256(ClientIP) AS hashed_ip          -- redact PII at ingest
FROM http_logs_stream
WHERE EdgeResponseStatus >= 400;

Configure via Dashboard (Logpush → Create a job → Basin Pipelines destination) or API.

Basin Pipelines + Queues Fan-out

await Promise.all([
  env.ANALYTICS_STREAM.send([event]),  // long-term storage + SQL
  env.PROCESS_QUEUE.send(event),       // immediate processing + retries
]);

Use Basin Pipelines for long-term storage + SQL; Queues for immediate processing/retries/DLQ; both for fan-out.

Observability (GraphQL Analytics)

Same R2 API token works. Endpoint: https://api.cloudflare.com/client/v4/graphql. Datasets cover ingestion, processing (incl. decodeErrors), delivery, sink writes (filesWritten), and user/validation errors — see the metrics doc for the full dataset/field catalog.

curl -X POST "https://api.cloudflare.com/client/v4/graphql" \
  -H "Authorization: Bearer $API_TOKEN" -H "Content-Type: application/json" \
  -d '{"query": "query { viewer { accounts(filter: {accountTag: \"'$ACCOUNT_ID'\"}) { pipelinesIngestionAdaptiveGroups(filter: {pipelineId: \"PIPELINE-UUID-WITH-DASHES\", datetime_geq: \"2026-03-01T00:00:00Z\"}, limit: 10) { sum { ingestedRecords ingestedBytes } dimensions { datetimeHour } } } } }"}'

Sink/pipeline IDs need dashes for GraphQL but wrangler may show them without: b909fe6e544844abbd63f6dcbc81d602 → b909fe6e-5448-44ab-bd63-f6dcbc81d602. Metrics take 5–10 min to populate.

Detecting Silent Data Loss

If a sink's bucket is deleted or its token expires, events are accepted but lost. Tell-tale: recordsWritten > 0 but filesWritten = 0. Always verify data lands in R2 within the roll interval and Basin SQL returns expected counts.

Schema Evolution (Immutable Basin Pipelines)

Basin Pipelines can't change. Version + dual-write:

npx wrangler basin pipelines streams create events_v2 --schema-file v2.json
await Promise.all([env.EVENTS_V1.send([event]), env.EVENTS_V2.send([event])]);
// query across versions with UNION ALL in Basin SQL

End-to-End: Streaming Analytics Dashboard

External APIs → Collector Worker (cron) → Pipeline → R2 (Iceberg) → Dashboard Worker → Basin SQL
  1. Create bucket + enable catalog (Basin Catalog)
  2. Create stream + sink + pipeline (here)
  3. Collector Worker with cron + stream binding (above)
  4. Dashboard Worker querying Basin SQL (sql/patterns.md)
  5. Enable automatic compaction

See Also

Source: SKILL.md on GitHub

No alertstoday3 checks · Risk SAFE
  • Gen Agent Trust Hubtoday

    This skill provides documentation and implementation patterns for Cloudflare Basin analytics workflows, including Pipelines, Catalog, and SQL querying. No security issues were detected, and the skill correctly leverages standard Cloudflare CLI tools and official API endpoints while following best practices for credential management.

  • Sockettoday

    No alerts

  • Snyktoday

    Risk: LOW · No issues

Signed by skilld at 41e0d19. This ties the file your Agent reads to that commit on GitHub. It does not review the instructions.

Last checked against GitHub 5 hours ago.

Activeupdated 6 hours ago

README badge

README badge for cloudflare/skills/basin