跳转到内容
搜索文档

写入 stream

最后更新 查看 MarkdownAgent 设置

通过 Worker 绑定(binding) 或 HTTP 端点向 stream 发送事件,适用于客户端应用和外部系统。

通过 Workers 发送

Worker 绑定提供了一种从 Workers 向 stream 发送数据的安全方式,无需管理 API 令牌或凭据。

配置 pipeline 绑定

在 Wrangler 文件中添加指向 stream 的 pipeline 绑定:

{
	"pipelines": [
		{
			"binding": "STREAM",
			"stream": "<STREAM_ID>"
		}
	]
}
[[pipelines]]
binding = "STREAM"
stream = "<STREAM_ID>"

Workers API

pipeline 绑定公开了一个用于向 stream 发送数据的方法:

send(records)

向 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>;

类型化 pipeline 绑定

当 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。

通过 HTTP 发送

每个 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 令牌文档。

Schema 验证

Stream 根据其配置以不同方式处理验证:

  • 结构化 stream:事件必须匹配定义的 schema 字段和类型。
  • 非结构化 stream:接受任何有效的 JSON 结构。数据存储在单个 value 列中。

对于结构化 stream,请确保事件与 schema 定义匹配。无效事件会被接受但会被丢弃,因此在发送前验证数据以避免事件丢失。使用 Worker 绑定时,运行 wrangler types 生成类型化 pipeline 绑定,在编译时捕获 schema 违规。您还可以查询用户错误指标以监控丢弃的事件并诊断 schema 验证问题。

这篇文档对您有帮助吗?