跳转到内容
搜索文档

入门指南

最后更新 查看 MarkdownAgent 设置

本指南将指导您完成以下操作:

  • 创建您的第一个 R2 bucket 并启用其 data catalog
  • 创建管道在数据目录进行身份验证所需的 API token
  • 使用简单的电子商务架构创建您的第一个管道,该管道可写入由 R2 Data Catalog 管理的 Apache Iceberg 表。
  • 通过 HTTP 端点发送示例电子商务数据。
  • 验证您的存储桶中的数据并使用 R2 SQL 查询它。

先决条件

  1. 注册 Cloudflare 账户
  2. 安装 Node.js

Node.js 版本管理器

使用 Voltanvm 等 Node 版本管理器,以避免权限问题并切换 Node.js 版本。本指南后续将介绍的 Wrangler 需要 Node 版本 16.17.0 或更高。

1. 创建 R2 存储桶

  1. 如果尚未登录,请运行:

    npx wrangler login
  2. 创建 R2 存储桶:

    npx wrangler r2 bucket create pipelines-tutorial
  1. 在 Cloudflare 仪表板中,转到 R2 object storage(R2 对象存储) 页面。

    Go to Overview ↗
  2. 选择 Create bucket(创建存储桶)

  3. 输入存储桶名称:pipelines-tutorial

  4. 选择 Create bucket(创建存储桶)

2. 启用 R2 Data Catalog

在您的 R2 存储桶上启用目录:

npx wrangler r2 bucket catalog enable pipelines-tutorial

当您运行此命令时,请记下 "Warehouse" 和 "Catalog URI"。稍后您将需要这些内容。

  1. 在 Cloudflare 仪表板中,转到 R2 object storage(R2 对象存储) 页面。

    Go to Overview ↗
  2. 选择存储桶:pipelines-tutorial。

  3. 切换到 Settings(设置) 选项卡,向下滚动到 R2 Data Catalog,然后选择 Enable(启用)

  4. 启用后,请记下 Catalog URIWarehouse name

3. 创建 API 令牌

管道必须使用具有目录和 R2 权限的 R2 API token 对 R2 Data Catalog 进行身份验证。

  1. 在 Cloudflare 仪表板中,转到 R2 object storage(R2 对象存储) 页面。

    Go to Overview ↗
  2. 选择 Manage API tokens(管理 API 令牌)

  3. 选择 Create Account API token(创建账户 API 令牌)

  4. 为您的 API 令牌命名。

  5. Permissions(权限) 下,选择 Admin Read & Write(管理员读写) 权限。

  6. 选择 Create Account API Token(创建账户 API 令牌)

  7. 记下 Token value

4. 创建管道

首先,创建一个定义您的电子商务数据结构的架构文件:

创建 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

按照提示操作:

  1. Pipeline name:输入 ecommerce

  2. Stream configuration

    • Enable HTTP endpoint: yes
    • Require authentication: no (为简单起见)
    • Configure custom CORS origins: no
    • Schema definition: Load from file
    • Schema file path: schema.json (或您的文件路径)
  3. 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 (为了在本教程中更快地看到数据)
  4. SQL transformation:选择 Use simple ingestion query 以使用:

    INSERT INTO ecommerce_sink SELECT * FROM ecommerce_stream

设置完成后,请记下最终输出中显示的 HTTP 端点 URL。

  1. 在 Cloudflare 仪表板中,转到 Pipelines(管道) > Pipelines(管道)

    Go to Pipelines ↗
  2. 选择 Create Pipeline(创建管道)

  3. Connect to a Stream(连接到 Stream)

    • Pipeline name(管道名称)ecommerce
    • Enable HTTP endpoint for sending data(启用用于发送数据的 HTTP 端点):已启用
    • HTTP authentication(HTTP 身份验证):已禁用(默认)
    • 选择 Next(下一步)
  4. 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(下一步)
  5. Define Sink(定义接收器)

    • 选择您的 R2 存储桶:pipelines-tutorial
    • Storage type(存储类型)R2 Data Catalog
    • Namespace(命名空间)default
    • Table name(表名称)ecommerce
    • Advanced Settings(高级设置):将 Maximum Time Interval(最大时间间隔) 更改为 10 seconds
    • 选择 Next(下一步)
  6. Credentials(凭据)

    • 禁用 Automatically create an Account API token for your sink(自动为接收器创建账户 API 令牌)
    • 输入第 3 步中的 Catalog Token(目录令牌)
    • 选择 Next(下一步)
  7. Pipeline Definition(管道定义)

    • 保留默认的 SQL 查询:
      INSERT INTO ecommerce_sink SELECT * FROM ecommerce_stream;
    • 选择 Create Pipeline(创建管道)
  8. 管道创建后,请记下 Stream ID(流 ID) 以备下一步使用。

5. 发送示例数据

将电子商务事件发送到您的管道 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} 替换为管道设置中的实际流端点。

6. 验证存储桶中的数据

  1. 在 Cloudflare 仪表板中,转到 R2 object storage(R2 对象存储) 页面。

  2. 选择您的存储桶:pipelines-tutorial

  3. 您应该会看到由您的管道创建的 Iceberg 元数据文件和数据文件。注意:如果您没有在您的存储桶中看到任何文件,请尝试等待几分钟并再次尝试。

  4. 数据按 Apache Iceberg 格式组织,其中包含用于跟踪表版本的元数据。

7. 使用 R2 SQL 查询您的数据

设置您的环境以使用 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

了解更多

管理 R2 Data Catalog

在您的存储桶上启用或禁用 R2 Data Catalog,检索配置详细信息,并对您的 Iceberg 引擎进行身份验证。

尝试另一个示例

有关设置简单的欺诈检测数据管道并在 Python 中为其生成事件的详细教程。

这篇文档对您有帮助吗?