跳转到内容
Telemetry
集成指南

Kafka 消费者延迟和处理分析

使用结构化事件监控 Kafka 消费者结果、分区延迟、处理延迟、重试和有害消息处理。

审阅者 Telemetry产品团队 . 埋点合约、隐私边界和实施指南. 审查标准和所有权

有用于
  • 消费者滞后监控
  • 分区不平衡
  • 消息处理可靠性
实施证据

Kafka Consumer Lag and Processing Analytics:从边界到验证行

在受控应用程序边界使用 Kafka Consumer Lag and Processing Analytics,保持事件契约较小,并在构建聚合视图之前验证已知结果。

  1. 1

    选择结果

    消费者滞后监控

  2. 2

    定义合同

    主题、分区、consumer_group 和状态

  3. 3

    埋点边界

    与高容量的每条消息完成事件分开发出聚合滞后快照。

  4. 4

    核实证据

    Exercise a known fixture, then inspect kafka_consumer_snapshot for one correctly typed terminal row.

开始之前

先决条件和界限

  • 服务器端TELEMETRY_API_KEY
  • 稳定的主题和消费者组名称
  • 来自 Kafka 客户端的偏移量和高水位元数据

交货设置

安装并初始化服务器端

在仅服务器代码中导入 telemetry-sh 并使用 process.env.TELEMETRY_API_KEY 对其进行一次初始化。 将摄取凭据保留在浏览器包、客户端可见的环境变量、源代码控制、日志和异常消息之外。

kafka-安装

npm安装

bash
npm install telemetry-sh
  1. 1准备一个具有有限网络行为的可重用服务器端交付客户端。
  2. 2在成功、失败、重试或超时边界处添加结果事件。
  3. 3在启用警报之前发送受控装置并检查存储的行。

片段

从一个结构化事件开始

在工作流程完成、失败或重试的位置添加此形状。然后从真实的字段构建仪表板。

kafka

Kafka Consumer Lag and Processing Analytics事件

javascript
await telemetry.log("kafka_consumer_snapshot", {
  topic,
  partition,
  consumer_group: groupId,
  offset: Number(currentOffset),
  high_watermark: Number(highWatermark),
  lag: Number(highWatermark - currentOffset),
  status: "healthy",
  consumer_instance: instanceId,
  release: process.env.APP_RELEASE,
});

活动合约

主题、分区、consumer_group 和状态

偏移、high_watermark、滞后和 duration_ms

尝试、error_type、consumer_instance 和释放

实施检查点

检查站 1

与高容量的每条消息完成事件分开发出聚合滞后快照。

检查站 2

请勿发送消息值、包含客户数据的密钥或身份验证配置。

检查站 3

保留主题和消费者组,同时将分区级图表限制为操作调查。

验证

证明事件已到达

在演练已知的成功和失败案例后运行此命令。如果您的最终事件契约与代码片段不同,请替换后备表名称。

kafka-验证

Kafka Consumer Lag and Processing Analytics 验证查询

sql
SELECT *
FROM kafka_consumer_snapshot
ORDER BY timestamp_utc DESC
LIMIT 20;
确认每个逻辑结果的一个终端行,以及预期状态、标识符、单位和 UTC 时间。
检查推断的架构并验证重试不会更改字段类型或生成新的逻辑事件 ID。
在存储的字段中搜索凭据、原始有效负载、提示、私有内容和无限制的错误消息。
在将仪表板视为完整之前,请执行提供程序超时、摄取拒绝和进程关闭。

实施参考

在启用新的生产路径之前,请检查事件合同、数据安全指南和上游主要文档。

生产边界

保持结果事件小且可恢复

该模式提供了

  • 除了上游工作流程之外,还有一个有界的、SQL 就绪的结果。
  • 用于仪表板、警报和跨事件关联的稳定字段。
  • 用于验证成功、失败、重试和超时行为的夹具驱动路径。

该模式不提供

  • OTLP 导出器、自动收集管道或详细跟踪和诊断日志的替代品。
  • 仅因为有效负载包含事件 ID,所以仅传送一次。
  • 收集原始提供商有效负载、用户内容、凭证或受监管数据的权限。

事件架构起点

在将查询或代码片段适应生产之前,请检查行粒度、发出边界、所需类型、隐私类、示例有效负载和验证清单。

相关产品功能

继续此工作流程 结构化事件

捕获稳定的事件名称、键入的字段和经过隐私审查的操作上下文。

相关 SQL 查询示例

用 SQL 回答下一个问题

针对此工作流程中的结构化字段运行查询,检查示例结果,并将有用的答案转换为仪表板或警报。

浏览所有查询示例

按实施系列浏览

比较相关集成模式

与此集成配对的模板

更多集成