
Cloudflare Pipelines 配置全指南從 Worker 綁定、Streams/Sinks 到不可變 SQL 管道的完整實戰手冊【免費下載鏈接】skillsSkills Catalog for Codex項目地址: https://gitcode.com/GitHub_Trending/skills4/skills本指南以 Cloudflare Deploy Skill 中 Pipelines 配置文檔 為骨架系統講解 Cloudflare PipelinesETL 流式平臺的完整配置鏈路在wrangler.jsonc中聲明 Worker 綁定、為結構化流定義 Schema、用 CLI 創建 Streams 與兩類 SinkIceberg 數據目錄 / Parquet 原始存儲、創建不可變 SQL 管道并落地到 R2。閱讀完成后你將能夠從零搭建一條數據源 → Streams → Pipelines(SQL) → Sinks → R2的生產級流式 ETL 管道并掌握 Schema 校驗、憑據配置、性能調參與常見故障排查的完整要點。一、Pipelines 是什么一個面向 R2 的流式 ETL 平臺Cloudflare Pipelines 是一個用于采集Ingest、轉換Transform并加載Load數據到 R2的流式 ETL 平臺。其核心架構由三部分組成見 Pipelines READMEData Sources → Streams → Pipelines (SQL) → Sinks → R2 ↑ ↓ ↓ HTTP/Workers Transform Iceberg/Parquet組件職責關鍵特性Streams事件采集持久的緩沖層HTTP / Workers 寫入支持結構化帶 Schema 校驗與非結構化兩種模式Pipelines用 SQL 對流做轉換創建后不可變immutable無法修改 SQLSinks將數據寫入 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——這是最容易踩的坑。通過以下命令獲取流 IDnpx wrangler pipelines streams listbinding是你自己在 Worker 代碼中使用的環境變量名如STREAM之后在Env類型中聲明并調用env.STREAM.send(...)即可寫入事件。修改綁定后必須重新部署npx wrangler deploy才會生效若出現env.STREAM is undefined優先檢查這兩點詳見 gotchas.md。綁定完成后Worker 中的最小寫入示例來自 api.mdinterface Env { STREAM: Pipeline; } export default { async fetch(request: Request, env: Env, ctx: ExecutionContext): PromiseResponse { 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 ExportedHandlerEnv;三、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 CatalogIceberg 表具備 ACID 事務、時間旅行time-travel、Schema 演化能力但配置更復雜需要 namespace、table、catalog token僅做文件存儲/歸檔→ 選R2 RawParquet/JSON 文件簡單直接但沒有內置 SQL 查詢配合外部工具Spark / Athena→ 選R2 RawParquet 分區標準格式 分區裁剪提升查詢性能但 Schema 兼容性需自行維護。5.1 R2 Data CatalogIcebergSinknpx 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 Catalognpx wrangler r2 bucket catalog enable my-bucket詳見 r2-data-catalog 配置文檔--catalog-token為具有R2 Admin Read Write權限的 API Token。5.2 R2 RawParquetSinknpx 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 關鍵參數速查表配置文檔原文OptionValuesGuidance--compressionzstd、snappy、gzipzstd壓縮比最佳snappy速度最快--roll-interval秒低延遲場景設 10–60查詢性能優先設 300--roll-sizeMB越大壓縮效果越好性能調優組合建議來自 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 purchase6.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 ... ENDCASE 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 tokenIceberg Sink 用R2 Admin Read WriteDashboard → R2 → API tokensR2 credentialsRaw Sink 用Object Read Writewrangler r2 bucket create的輸出HTTP ingest token外部寫入用Workers Pipeline SendDashboard → 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_TOKENYOUR_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 loginCI/CD 用CLOUDFLARE_API_TOKEN環境變量。十、調試清單與常見錯誤診斷清單來自 gotchas.md流存在npx wrangler pipelines streams list管道健康npx wrangler pipelines get IDSQL 語法與 Schema 字段匹配添加綁定后已重新部署 Worker已等待 roll interval10–300 秒Accepted 數量與 Processed 數量一致無校驗靜默丟棄常見錯誤對照錯誤原因修復事件不在 R2 中roll interval 未到等待 10–300s檢查roll_intervalSchema 校驗失敗類型不匹配、缺必填字段客戶端先行校驗限流429單流寫入 5 MB/s批量發送、申請提高額度負載過大413單請求 1 MB拆分為更小的批次無法刪除流仍有 Pipeline 引用先刪除管道Sink 憑據錯誤Token 過期用新憑據重建 SinkOpen 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),僅供參考