Cloudflare Pipelines 配置全指南:从 Worker 绑定、Streams/Sinks 到不可变 SQL 管道的完整实战手册
【免费下载链接】skillsSkills Catalog for Codex项目地址: https://gitcode.com/GitHub_Trending/skills4/skills
本指南以 Cloudflare Deploy Skill 中 Pipelines 配置文档 为骨架,系统讲解 Cloudflare Pipelines(ETL 流式平台)的完整配置链路:在wrangler.jsonc中声明 Worker 绑定、为结构化流定义 Schema、用 CLI 创建 Streams 与两类 Sink(Iceberg 数据目录 / Parquet 原始存储)、创建不可变 SQL 管道并落地到 R2。阅读完成后,你将能够从零搭建一条"数据源 → Streams → Pipelines(SQL) → Sinks → R2"的生产级流式 ETL 管道,并掌握 Schema 校验、凭据配置、性能调参与常见故障排查的完整要点。
一、Pipelines 是什么:一个面向 R2 的流式 ETL 平台
Cloudflare Pipelines 是一个用于采集(Ingest)、转换(Transform)并加载(Load)数据到 R2的流式 ETL 平台。其核心架构由三部分组成(见 Pipelines README):
Data Sources → Streams → Pipelines (SQL) → Sinks → R2 ↑ ↓ ↓ HTTP/Workers Transform Iceberg/Parquet| 组件 | 职责 | 关键特性 |
|---|---|---|
| Streams | 事件采集,持久的缓冲层(HTTP / Workers 写入) | 支持结构化(带 Schema 校验)与非结构化两种模式 |
| Pipelines | 用 SQL 对流做转换 | 创建后不可变(immutable),无法修改 SQL |
| Sinks | 将数据写入 R2 目的地 | 提供精确一次(exactly-once)投递语义 |
常见的落地形态:分析管道(点击流、遥测、服务器日志)、数据仓库(ETL 进可查询的 Iceberg 表)、事件处理(移动端 / IoT 富化)、电商分析(用户事件、购买、浏览)。当前状态为Open Beta,需要 Workers Paid 套餐;除标准 R2 存储/操作费用外不额外计费(以仓库文档声明为准)。
二、Worker Binding:在 wrangler.jsonc 中声明流绑定
要让 Worker 代码能够写入流,首先需要在wrangler.jsonc中声明Pipelines 绑定。配置文档给出的完整示例:
// wrangler.jsonc { "pipelines": [ { "pipeline": "<STREAM_ID>", "binding": "STREAM" } ] }要点解析:
pipeline字段填的是Stream ID(流 ID),而不是 Pipeline(管道)ID——这是最容易踩的坑。通过以下命令获取流 ID:npx wrangler pipelines streams listbinding是你自己在 Worker 代码中使用的环境变量名(如STREAM),之后在Env类型中声明并调用env.STREAM.send(...)即可写入事件。- 修改绑定后必须重新部署(
npx wrangler deploy)才会生效;若出现env.STREAM is undefined,优先检查这两点(详见 gotchas.md)。
绑定完成后,Worker 中的最小写入示例(来自 api.md):
interface Env { STREAM: Pipeline; } export default { async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> { const event = { user_id: "123", event_type: "purchase", amount: 29.99 }; // Fire-and-forget 模式:不阻塞响应 ctx.waitUntil(env.STREAM.send([event])); return new Response('OK'); } } satisfies ExportedHandler<Env>;三、Schema:为结构化流定义字段与类型
创建结构化流时,可以用一个 JSON 文件描述事件结构,Pipelines 会据此在写入时做字段级校验。配置文档中的标准 Schema 示例:
{ "fields": [ { "name": "user_id", "type": "string", "required": true }, { "name": "event_type", "type": "string", "required": true }, { "name": "amount", "type": "float64", "required": false }, { "name": "timestamp", "type": "timestamp", "required": true } ] }每个字段由三个属性组成:
| 属性 | 含义 | 说明 |
|---|---|---|
name | 字段名 | 与 SQL 转换中引用的列名一致 |
type | 字段类型 | 见下方支持类型清单 |
required | 是否必填 | true缺失即校验失败;false允许缺省 |
支持的类型(配置文档原文):
string、int32、int64、float32、float64、bool、timestamp、json、binary、list、struct
⚠️ 重要提示(Gotcha):结构化流对不合法的事件会静默丢弃——HTTP 返回 200 但事件永远不会出现在 Sink 中(详见 gotchas.md)。因此强烈建议在客户端用 Zod 先做一次校验,获得即时反馈:
import { z } from 'zod'; const EventSchema = z.object({ user_id: z.string(), event_type: z.enum(['purchase', 'view']), amount: z.number().positive().optional() }); try { const validated = EventSchema.parse(rawEvent); // 校验失败会抛异常 await env.STREAM.send([validated]); } catch (e) { // 在此获得即时错误反馈 }四、Stream Setup:创建、查询与删除流
流的创建有两种模式:
# 带 Schema(结构化流,写入时校验) npx wrangler pipelines streams create my-stream --schema-file schema.json # 不带 Schema(非结构化流,无校验) npx wrangler pipelines streams create my-stream生命周期管理命令:
# 列出所有流(拿到 Stream ID) npx wrangler pipelines streams list # 查看单个流详情 npx wrangler pipelines streams get <ID> # 删除流(注意:若还有 Pipeline 引用它,需先删除管道) npx wrangler pipelines streams delete <ID>实战提示:若删除流时提示失败,通常是因为该流仍被某个 Pipeline 引用,应先删除管道再删流(见 gotchas.md 错误对照表)。
向流写入事件的两种途径:
- Worker 绑定(上文已述):
env.STREAM.send(events),支持单个对象或数组;单次请求上限1 MB,单流写入速率上限5 MB/s。 - HTTP Ingest 端点(适用于外部系统/服务端上报):
curl -X POST https://{stream-id}.ingest.cloudflare.com \ -H "Content-Type: application/json" \ -H "Authorization: Bearer YOUR_API_TOKEN" \ -d '[{"user_id": "123", "event_type": "purchase"}]'其中{stream-id}同样来自npx wrangler pipelines streams list,鉴权 Token 需要Workers Pipeline Send权限(Dashboard → Workers → API tokens)。关键细节:HTTP 端点的请求体必须是 JSON 数组,而不是单个对象,否则会返回 400(详见 api.md)。
五、Sink Configuration:两类 R2 落地目标
Sink 决定数据最终以什么形态写入 R2。配置文档给出了两种类型,选择依据如下(来自 README):
- 需要直接对数据跑 SQL 查询→ 选R2 Data Catalog(Iceberg 表):具备 ACID 事务、时间旅行(time-travel)、Schema 演化能力,但配置更复杂(需要 namespace、table、catalog token);
- 仅做文件存储/归档→ 选R2 Raw(Parquet/JSON 文件):简单直接,但没有内置 SQL 查询;
- 配合外部工具(Spark / Athena)→ 选R2 Raw(Parquet + 分区):标准格式 + 分区裁剪提升查询性能,但 Schema 兼容性需自行维护。
5.1 R2 Data Catalog(Iceberg)Sink
npx wrangler pipelines sinks create my-sink \ --type r2-data-catalog \ --bucket my-bucket --namespace default --table events \ --catalog-token $TOKEN \ --compression zstd --roll-interval 60其中--bucket指定的桶必须先启用 Data Catalog(npx wrangler r2 bucket catalog enable my-bucket,详见 r2-data-catalog 配置文档),--catalog-token为具有R2 Admin Read & Write权限的 API Token。
5.2 R2 Raw(Parquet)Sink
npx wrangler pipelines sinks create my-sink \ --type r2 --bucket my-bucket --format parquet \ --path analytics/events \ --partitioning "year=%Y/month=%m/day=%d" \ --access-key-id $KEY --secret-access-key $SECRET--partitioning使用 strftime 风格的占位符做 Hive 风格分区目录,方便外部查询引擎做分区裁剪。
5.3 关键参数速查表(配置文档原文)
| Option | Values | Guidance |
|---|---|---|
--compression | zstd、snappy、gzip | zstd压缩比最佳,snappy速度最快 |
--roll-interval | 秒 | 低延迟场景设 10–60;查询性能优先设 300 |
--roll-size | MB | 越大压缩效果越好 |
性能调优组合建议(来自 patterns.md):
| 目标 | 配置 |
|---|---|
| 低延迟 | --roll-interval 10 |
| 查询性能 | --roll-interval 300 --roll-size 100 |
| 成本最优 | --compression zstd --roll-interval 300 |
六、Pipeline Creation:创建不可变 SQL 管道
Pipeline 负责用 SQL 把 Stream 中的数据转换后写入 Sink,命令格式为:
npx wrangler pipelines create my-pipeline \ --sql "INSERT INTO my_sink SELECT * FROM my_stream WHERE event_type = 'purchase'"6.1 常用 SQL 转换模式
管道中的 SQL 遵循INSERT INTO sink SELECT ... FROM stream结构,常用模式包括(详见 api.md 与 patterns.md):
-- ① 过滤事件(尽早裁剪,减少存储) INSERT INTO my_sink SELECT * FROM my_stream WHERE event_type = 'purchase' AND amount > 100 -- ② 只选需要的字段 INSERT INTO my_sink SELECT user_id, event_type, timestamp, amount FROM my_stream -- ③ 转换与富化(UPPER / 数学运算 / CONCAT / CASE WHEN) INSERT INTO my_sink SELECT user_id, UPPER(event_type) as event_type, timestamp, amount * 1.1 as amount_with_tax, CONCAT(user_id, '_', product_id) as unique_key, CASE WHEN amount > 1000 THEN 'high_value' WHEN amount > 100 THEN 'medium_value' ELSE 'low_value' END as customer_tier FROM my_stream WHERE event_type IN ('purchase', 'refund')6.2 可用 SQL 函数速查
| 函数 | 示例 | 用途 |
|---|---|---|
UPPER(s) | UPPER(event_type) | 字符串归一化 |
LOWER(s) | LOWER(email) | 大小写不敏感匹配 |
CONCAT(...) | CONCAT(user_id, '_', product_id) | 生成复合键 |
CASE WHEN ... THEN ... END | CASE WHEN amount > 100 THEN 'high' ELSE 'low' END | 条件富化 |
CAST(x AS type) | CAST(timestamp AS string) | 类型转换 |
COALESCE(x, y) | COALESCE(amount, 0.0) | 默认值兜底 |
| 数学运算符 | amount * 1.1、price / quantity | 计算 |
| 比较运算 | amount > 100、status IN ('active', 'pending') | 过滤 |
CAST支持的字符串类型与 Schema 类型一致:string、int32、int64、float32、float64、bool、timestamp。
6.3 ⚠️ Pipelines 是不可变的
创建后无法修改 SQL,只能删除重建:
npx wrangler pipelines delete old-pipeline npx wrangler pipelines create new-pipeline --sql "..."配套的 SQL 限制包括:不支持 JOIN(单管道只处理单个流)、不支持窗口函数、不支持子查询、无 Schema 演化(详见 gotchas.md)。因此官方建议:
- 使用版本化命名(如
events-pipeline-v1); - 将 SQL纳入版本控制;
- Schema 演化时采用"双写过渡"策略:创建 v2 流/管道后,
await Promise.all([env.EVENTS_V1.send([event]), env.EVENTS_V2.send([event])])同时写入新旧版本,过渡期结束后删除旧管道(详见 patterns.md)。
七、Credentials:三类凭据速查
| 类型 | 所需权限 | 获取位置 |
|---|---|---|
| Catalog token(Iceberg Sink 用) | R2 Admin Read & Write | Dashboard → R2 → API tokens |
| R2 credentials(Raw Sink 用) | Object Read & Write | wrangler r2 bucket create的输出 |
| HTTP ingest token(外部写入用) | Workers Pipeline Send | Dashboard → Workers → API tokens |
安全建议(来自 r2-data-catalog 配置文档):Token 应通过环境变量或密钥管理器保存、绝不硬编码进代码;遵循最小权限原则(查询引擎用只读 Token,写入才用读写 Token);定期轮换 Token;每个应用单独建 Token 以便追踪与吊销。
八、查询落库数据(R2 Data Catalog)
若 Sink 是 Iceberg 表,可用 Wrangler 直接对数据跑标准 SQL(含 GROUP BY、JOIN、WHERE、ORDER BY 等):
export WRANGLER_R2_SQL_AUTH_TOKEN=YOUR_CATALOG_TOKEN npx wrangler r2 sql query "warehouse_name" " SELECT event_type, COUNT(*) as event_count, SUM(amount) as total_revenue FROM default.my_table WHERE event_type = 'purchase' AND timestamp >= '2025-01-01' GROUP BY event_type ORDER BY total_revenue DESC LIMIT 100"九、完整示例:一条生产级 ETL 管道的落地全流程
配置文档给出的端到端流程(my-bucket需为启用了 Data Catalog 的桶):
# 1. 创建并启用 R2 桶的 Data Catalog npx wrangler r2 bucket create my-bucket npx wrangler r2 bucket catalog enable my-bucket # 2. 用 Schema 文件创建结构化流 npx wrangler pipelines streams create my-stream --schema-file schema.json # 3. 创建 Iceberg Sink(按需调整 --catalog-token 等参数) npx wrangler pipelines sinks create my-sink --type r2-data-catalog --bucket my-bucket ... # 4. 创建不可变 SQL 管道 npx wrangler pipelines create my-pipeline --sql "INSERT INTO my_sink SELECT * FROM my_stream" # 5. 部署 Worker(让绑定生效) npx wrangler deploy部署前请确认已通过npx wrangler whoami完成认证(未认证时参考 wrangler/auth.md:本地用wrangler login,CI/CD 用CLOUDFLARE_API_TOKEN环境变量)。
十、调试清单与常见错误
诊断清单(来自 gotchas.md):
- 流存在:
npx wrangler pipelines streams list - 管道健康:
npx wrangler pipelines get <ID> - SQL 语法与 Schema 字段匹配
- 添加绑定后已重新部署 Worker
- 已等待 roll interval(10–300 秒)
- Accepted 数量与 Processed 数量一致(无校验静默丢弃)
常见错误对照:
| 错误 | 原因 | 修复 |
|---|---|---|
| 事件不在 R2 中 | roll interval 未到 | 等待 10–300s,检查roll_interval |
| Schema 校验失败 | 类型不匹配、缺必填字段 | 客户端先行校验 |
| 限流(429) | 单流写入 >5 MB/s | 批量发送、申请提高额度 |
| 负载过大(413) | 单请求 >1 MB | 拆分为更小的批次 |
| 无法删除流 | 仍有 Pipeline 引用 | 先删除管道 |
| Sink 凭据错误 | Token 过期 | 用新凭据重建 Sink |
Open Beta 限制(以仓库文档为准):每个账户 Streams/Sinks/Pipelines 各 20 个;Payload 上限 1 MB;单流摄入速率 5 MB/s;事件保留 24 小时;推荐批量大小 100 个事件。
延伸阅读
本指南聚焦配置链路,若需继续深入,建议按 Pipelines README 的阅读顺序展开:
- api.md —— 发送事件、TypeScript 类型、SQL 函数全参考、HTTP 响应码
- patterns.md —— Fire-and-forget、Zod 校验、Pipelines + Queues 扇出、性能调优、Schema 版本化
- gotchas.md —— 静默丢弃、不可变管道、限流与限制
- r2-data-catalog —— 桶的 Catalog 启用、PyIceberg 客户端配置与 Token 权限细节
- r2 —— R2 桶管理、S3 SDK 接入、生命周期与事件通知
【免费下载链接】skillsSkills Catalog for Codex项目地址: https://gitcode.com/GitHub_Trending/skills4/skills
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考