Kafka Consumer Lag and Processing Analytics:從邊界到驗證行
在受控應用程式邊界使用 Kafka Consumer Lag and Processing Analytics,保持事件契約較小,並在建置聚合檢視之前驗證已知結果。
- 1
選擇結果
消費者滯後監控
- 2
定義合約
主題、分割區、consumer_group 和狀態
- 3
埋點邊界
與高容量的每條訊息完成事件分開發出聚合滯後快照。
- 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 對其進行一次初始化。 將攝取憑據保留在瀏覽器包、客戶端可見的環境變數、原始碼控制、日誌和異常訊息之外。
npm安裝
npm install telemetry-sh- 1準備一個具有有限網路行為的可重用伺服器端交付客戶端。
- 2在成功、失敗、重試或超時邊界處新增結果事件。
- 3在啟用警示之前傳送受控裝置並檢查儲存的行。
片段
從一個結構化事件開始
在工作流程完成、失敗或重試的位置新增此形狀。然後從真實的欄位建置儀表板。
Kafka Consumer Lag and Processing Analytics事件
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 Consumer Lag and Processing Analytics 驗證查詢
SELECT *
FROM kafka_consumer_snapshot
ORDER BY timestamp_utc DESC
LIMIT 20;實施參考
在啟用新的生產路徑之前,請檢查事件合約、資料安全指南和上游主要文件。
生產邊界
保持結果事件小且可恢復
該模式提供了
- 除了上游工作流程之外,還有一個有界的、SQL 就緒的結果。
- 用於儀表板、警示和跨事件關聯的穩定欄位。
- 用於驗證成功、失敗、重試和超時行為的夾具驅動路徑。
該模式不提供
- OTLP 匯出器、自動收集管道或詳細追蹤和診斷日誌的替代品。
- 僅因為有效負載包含事件 ID,所以僅傳送一次。
- 收集原始提供商有效負載、使用者內容、憑證或受監管資料的權限。
事件架構起點
此工作流程的事件契約
在將查詢或程式碼片段適應生產之前,請檢查行粒度、發出邊界、所需型別、隱私類、範例有效負載和驗證清單。
相關產品功能
繼續此工作流程 結構化事件
捕獲穩定的事件名稱、鍵入的欄位和經過隱私審查的操作上下文。
相關 SQL 查詢範例
用 SQL 回答下一個問題
針對此工作流程中的結構化欄位執行查詢,檢查範例結果,並將有用的答案轉換為儀表板或警示。
按實施系列瀏覽
比較相關整合模式
與此整合配對的範本
更多整合