本指南将指导您完成以下操作:
- 创建您的第一个 R2 bucket 并启用其 data catalog。
- 创建管道在数据目录进行身份验证所需的 API token。
- 使用简单的电子商务架构创建您的第一个管道,该管道可写入由 R2 Data Catalog 管理的 Apache Iceberg ↗ 表。
- 通过 HTTP 端点发送示例电子商务数据。
- 验证您的存储桶中的数据并使用 R2 SQL 查询它。
- 注册 Cloudflare 账户 ↗。
- 安装
Node.js↗。
Node.js 版本管理器
使用 Volta ↗ 或 nvm ↗ 等 Node 版本管理器,以避免权限问题并切换 Node.js 版本。本指南后续将介绍的 Wrangler 需要 Node 版本 16.17.0 或更高。
-
如果尚未登录,请运行:
npx wrangler login -
创建 R2 存储桶:
npx wrangler r2 bucket create pipelines-tutorial
-
在 Cloudflare 仪表板中,转到 R2 object storage(R2 对象存储) 页面。
Go to Overview ↗ -
选择 Create bucket(创建存储桶)。
-
输入存储桶名称:pipelines-tutorial
-
选择 Create bucket(创建存储桶)。
在您的 R2 存储桶上启用目录:
npx wrangler r2 bucket catalog enable pipelines-tutorial当您运行此命令时,请记下 "Warehouse" 和 "Catalog URI"。稍后您将需要这些内容。
-
在 Cloudflare 仪表板中,转到 R2 object storage(R2 对象存储) 页面。
Go to Overview ↗ -
选择存储桶:pipelines-tutorial。
-
切换到 Settings(设置) 选项卡,向下滚动到 R2 Data Catalog,然后选择 Enable(启用)。
-
启用后,请记下 Catalog URI 和 Warehouse name。
管道必须使用具有目录和 R2 权限的 R2 API token 对 R2 Data Catalog 进行身份验证。
-
在 Cloudflare 仪表板中,转到 R2 object storage(R2 对象存储) 页面。
Go to Overview ↗ -
选择 Manage API tokens(管理 API 令牌)。
-
选择 Create Account API token(创建账户 API 令牌)。
-
为您的 API 令牌命名。
-
在 Permissions(权限) 下,选择 Admin Read & Write(管理员读写) 权限。
-
选择 Create Account API Token(创建账户 API 令牌)。
-
记下 Token value。
首先,创建一个定义您的电子商务数据结构的架构文件:
创建 schema.json:
{
"fields": [
{
"name": "user_id",
"type": "string",
"required": true
},
{
"name": "event_type",
"type": "string",
"required": true
},
{
"name": "product_id",
"type": "string",
"required": false
},
{
"name": "amount",
"type": "float64",
"required": false
}
]
}使用交互式设置创建可写入 R2 Data Catalog 的管道:
npx wrangler pipelines setup按照提示操作:
-
Pipeline name:输入
ecommerce -
Stream configuration:
- Enable HTTP endpoint:
yes - Require authentication:
no(为简单起见) - Configure custom CORS origins:
no - Schema definition:
Load from file - Schema file path:
schema.json(或您的文件路径)
- Enable HTTP endpoint:
-
Sink configuration:
- Destination type:
Data Catalog Table - R2 bucket name:
pipelines-tutorial - Namespace:
default - Table name:
ecommerce - Catalog API token: 输入第 3 步中的令牌
- Compression:
zstd - Roll file when size reaches (MB):
100 - Roll file when time reaches (seconds):
10(为了在本教程中更快地看到数据)
- Destination type:
-
SQL transformation:选择
Use simple ingestion query以使用:INSERT INTO ecommerce_sink SELECT * FROM ecommerce_stream
设置完成后,请记下最终输出中显示的 HTTP 端点 URL。
-
在 Cloudflare 仪表板中,转到 Pipelines(管道) > Pipelines(管道)。
Go to Pipelines ↗ -
选择 Create Pipeline(创建管道)。
-
Connect to a Stream(连接到 Stream):
- Pipeline name(管道名称):
ecommerce - Enable HTTP endpoint for sending data(启用用于发送数据的 HTTP 端点):已启用
- HTTP authentication(HTTP 身份验证):已禁用(默认)
- 选择 Next(下一步)
- Pipeline name(管道名称):
-
Define Input Schema(定义输入架构):
- 选择 JSON editor(JSON 编辑器)
- 复制架构:
{ "fields": [ { "name": "user_id", "type": "string", "required": true }, { "name": "event_type", "type": "string", "required": true }, { "name": "product_id", "type": "string", "required": false }, { "name": "amount", "type": "f64", "required": false } ] } - 选择 Next(下一步)
-
Define Sink(定义接收器):
- 选择您的 R2 存储桶:
pipelines-tutorial - Storage type(存储类型):R2 Data Catalog
- Namespace(命名空间):
default - Table name(表名称):
ecommerce - Advanced Settings(高级设置):将 Maximum Time Interval(最大时间间隔) 更改为
10 seconds - 选择 Next(下一步)
- 选择您的 R2 存储桶:
-
Credentials(凭据):
- 禁用 Automatically create an Account API token for your sink(自动为接收器创建账户 API 令牌)
- 输入第 3 步中的 Catalog Token(目录令牌)
- 选择 Next(下一步)
-
Pipeline Definition(管道定义):
- 保留默认的 SQL 查询:
INSERT INTO ecommerce_sink SELECT * FROM ecommerce_stream; - 选择 Create Pipeline(创建管道)
- 保留默认的 SQL 查询:
-
管道创建后,请记下 Stream ID(流 ID) 以备下一步使用。
将电子商务事件发送到您的管道 HTTP 端点:
curl -X POST https://{stream-id}.ingest.cloudflare.com \
-H "Content-Type: application/json" \
-d '[
{
"user_id": "user_12345",
"event_type": "purchase",
"product_id": "widget-001",
"amount": 29.99
},
{
"user_id": "user_67890",
"event_type": "view_product",
"product_id": "widget-002"
},
{
"user_id": "user_12345",
"event_type": "add_to_cart",
"product_id": "widget-003",
"amount": 15.50
}
]'将 {stream-id} 替换为管道设置中的实际流端点。
-
在 Cloudflare 仪表板中,转到 R2 object storage(R2 对象存储) 页面。
-
选择您的存储桶:
pipelines-tutorial。 -
您应该会看到由您的管道创建的 Iceberg 元数据文件和数据文件。注意:如果您没有在您的存储桶中看到任何文件,请尝试等待几分钟并再次尝试。
-
数据按 Apache Iceberg 格式组织,其中包含用于跟踪表版本的元数据。
设置您的环境以使用 R2 SQL:
export WRANGLER_R2_SQL_AUTH_TOKEN=YOUR_API_TOKEN或者创建一个包含以下内容的 .env 文件:
WRANGLER_R2_SQL_AUTH_TOKEN=YOUR_API_TOKEN其中 YOUR_API_TOKEN 是您在第 3 步中创建的令牌。有关设置环境变量的更多信息,请参阅 Wrangler system environment variables。
查询您的数据:
npx wrangler r2 sql query "YOUR_WAREHOUSE_NAME" "
SELECT
user_id,
event_type,
product_id,
amount
FROM default.ecommerce
WHERE event_type = 'purchase'
LIMIT 10"将 YOUR_WAREHOUSE_NAME 替换为第 2 步中的仓库名称。
您还可以使用任何支持 Apache Iceberg 的引擎查询该表。要了解有关将其他引擎连接到 R2 Data Catalog 的更多信息,请参阅 Connect to Iceberg engines。