通过 Worker 绑定(binding) 或 HTTP 端点向 stream 发送事件,适用于客户端应用和外部系统。
Worker 绑定提供了一种从 Workers 向 stream 发送数据的安全方式,无需管理 API 令牌或凭据。
在 Wrangler 文件中添加指向 stream 的 pipeline 绑定:
{
"pipelines": [
{
"binding": "STREAM",
"stream": "<STREAM_ID>"
}
]
}[[pipelines]]
binding = "STREAM"
stream = "<STREAM_ID>"pipeline 绑定公开了一个用于向 stream 发送数据的方法:
向 stream 发送 JSON 可序列化记录数组。返回一个在记录被确认摄取后 resolve 的 Promise。
export default {
async fetch(request, env, ctx) {
const events = await request.json();
await env.STREAM.send(events);
return new Response("Events sent");
},
};export default {
async fetch(request, env, ctx): Promise<Response> {
const events = await request.json<Record<string, unknown>[]>();
await env.STREAM.send(events);
return new Response("Events sent");
},
} satisfies ExportedHandler<Env>;当 stream 定义了 schema 时,运行 wrangler types 会为 pipeline 绑定生成特定于 schema 的 TypeScript 类型。绑定将获得带有完整自动补全和编译时类型检查的命名记录类型,而不是通用的 Pipeline<PipelineRecord>。有关更多信息,请参阅 wrangler types 文档。
运行 wrangler types 后,生成的 worker-configuration.d.ts 文件在 Cloudflare 命名空间内包含一个命名记录类型。类型名称派生自 stream 名称(而非绑定名称),转换为 PascalCase 并加上 Record 后缀。
以下是名为 ecommerce_stream 的 stream 在 worker-configuration.d.ts 中生成的类型示例:
declare namespace Cloudflare {
type EcommerceStreamRecord = {
user_id: string;
event_type: string;
product_id?: string;
amount?: number;
};
interface Env {
STREAM: import("cloudflare:pipelines").Pipeline<Cloudflare.EcommerceStreamRecord>;
}
}wrangler types 在以下情况下回退到通用的 Pipeline<PipelineRecord> 类型:
- 未认证:运行
wrangler login以启用类型化 pipeline 绑定。 - 未找到 stream:Wrangler 配置中的 stream ID 与现有 stream 不匹配。
- 非结构化 stream:stream 创建时未定义 schema。
每个 stream 都提供一个可选的 HTTP 端点,用于从外部应用、浏览器或任何能够发起 HTTP 请求的系统摄取数据。
HTTP 端点遵循以下格式:
https://{stream-id}.ingest.cloudflare.com在 Cloudflare 仪表板的 Pipelines > Streams 下查找 stream 的端点 URL,或使用 Wrangler CLI 通过 stream ID 或 stream 名称查询:
npx wrangler pipelines streams get <STREAM_NAME_OR_ID>通过 POST 请求以 JSON 数组形式发送事件:
curl -X POST https://{stream-id}.ingest.cloudflare.com \
-H "Content-Type: application/json" \
-d '[
{
"user_id": "12345",
"event_type": "purchase",
"product_id": "widget-001",
"amount": 29.99
}
]'当 stream 启用了身份验证时,在 Authorization 标头中包含 API 令牌:
curl -X POST https://{stream-id}.ingest.cloudflare.com \
-H "Content-Type: application/json" \
-H "Authorization: Bearer YOUR_API_TOKEN" \
-d '[{"event": "test"}]'API 令牌必须具有 Workers Pipeline Send 权限。有关更多信息,请参阅创建 API 令牌文档。
Stream 根据其配置以不同方式处理验证:
- 结构化 stream:事件必须匹配定义的 schema 字段和类型。
- 非结构化 stream:接受任何有效的 JSON 结构。数据存储在单个
value列中。
对于结构化 stream,请确保事件与 schema 定义匹配。无效事件会被接受但会被丢弃,因此在发送前验证数据以避免事件丢失。使用 Worker 绑定时,运行 wrangler types 生成类型化 pipeline 绑定,在编译时捕获 schema 违规。您还可以查询用户错误指标以监控丢弃的事件并诊断 schema 验证问题。