了解如何:
- 创建和配置用于数据摄取的 stream
- 查看和更新 stream 设置
- 在不再需要时删除 stream
Stream 通过 stream 名称作为 SQL 表提供给 pipeline(例如 SELECT * FROM my_stream)。
-
在 Cloudflare 仪表板中,前往 Pipelines 页面。
Go to Pipelines ↗ -
选择 Create Pipeline(创建 Pipeline) 启动 pipeline 创建向导。
-
完成向导以创建 stream 以及关联的 sink 和 pipeline。
要创建 stream,请运行 pipelines streams create 命令:
npx wrangler pipelines streams create <STREAM_NAME>或者,要运行帮助您配置 stream、sink 和 pipeline 的交互式设置向导,请运行 pipelines setup 命令:
npx wrangler pipelines setupStream 支持两种数据处理方式:
- 结构化 stream:定义具有特定字段和数据类型的 schema。事件将根据 schema 进行验证。
- 非结构化 stream:接受任何有效的 JSON 而不进行验证。这些 stream 有一个包含 JSON 数据的
value列。
要创建结构化 stream,请提供 schema 文件:
npx wrangler pipelines streams create my-stream --schema-file schema.jsonSchema 文件示例:
{
"fields": [
{
"name": "user_id",
"type": "string",
"required": true
},
{
"name": "amount",
"type": "float64",
"required": false
},
{
"name": "tags",
"type": "list",
"required": false,
"items": {
"type": "string"
}
},
{
"name": "metadata",
"type": "struct",
"required": false,
"fields": [
{
"name": "source",
"type": "string",
"required": false
},
{
"name": "priority",
"type": "int32",
"required": false
}
]
}
]
}支持的数据类型:
string- 文本值int32、int64- 整数float32、float64- 浮点数bool- 布尔值 true/falsetimestamp- RFC 3339 时间戳,或解析为 Unix 秒、毫秒或微秒的数字值(取决于单位)json- JSON 对象binary- 二进制数据(base64 编码)list- 值数组struct- 具有定义字段的嵌套对象
-
在 Cloudflare 仪表板中,前往 Pipelines > Streams(流)。
-
选择一个 stream 以查看其关联配置。
要查看特定 stream,请使用 stream ID 或 stream 名称运行 pipelines streams get 命令:
npx wrangler pipelines streams get <STREAM_NAME_OR_ID>要列出账户中的所有 stream,请运行 pipelines streams list 命令:
npx wrangler pipelines streams liststream 创建后可以更新某些 HTTP 摄取设置。stream 创建后不支持修改 schema。
-
在 Cloudflare 仪表板中,前往 Pipelines > Streams(流)。
-
选择要更新的 stream。
-
在 Settings(设置) 选项卡中,前往 HTTP Ingest(HTTP 摄取)。
-
要启用或禁用 HTTP 摄取,请选择 Enable(启用) 或 Disable(禁用)。
-
要更新身份验证和 CORS 设置,请选择 Edit(编辑) 并进行修改。
-
保存更改。
-
在 Cloudflare 仪表板中,前往 Pipelines > Streams(流)。
-
选择要删除的 stream。
-
在 Settings(设置) 选项卡中,前往 General(常规),然后选择 Delete(删除)。
要删除 stream,请运行 pipelines streams delete 命令:
npx wrangler pipelines streams delete <STREAM_ID>