You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何基于Azure Log Analytics的Databricks日志表编写KQL获取流处理指标

Databricks 流处理指标 KQL 查询示例

以下查询均基于你已同步的三类 Spark 日志表实现,可直接在 Azure Log Analytics 查询编辑器中运行,支持按需调整时间范围、集群名称等过滤条件。


1. 每秒处理的消息数量

基于 SparkMetrics_CL 中流处理完成指标聚合统计:

SparkMetrics_CL
| where Name_s == "streaming.processedRecords.perSecond"
| extend ClusterName = ClusterName_s, ProcessedPerSec = Value_d
| project TimeGenerated, ClusterName, ProcessedPerSec
// 按需过滤指定集群:| where ClusterName == "你的集群名称"
| render timechart with (title="每秒处理消息数量趋势")

2. 每秒输入、输出消息数量

分别拉取输入、输出速率指标合并展示:

let inputRate = SparkMetrics_CL
| where Name_s == "streaming.inputRecords.perSecond"
| project TimeGenerated, ClusterName = ClusterName_s, InputPerSec = Value_d;
let outputRate = SparkMetrics_CL
| where Name_s == "streaming.outputRecords.perSecond"
| project TimeGenerated, ClusterName = ClusterName_s, OutputPerSec = Value_d;
inputRate
| join kind=inner outputRate on TimeGenerated, ClusterName
| project TimeGenerated, ClusterName, InputPerSec, OutputPerSec
| render timechart with (title="每秒输入/输出消息数量")

3. 消息延迟

取流处理批次处理延迟指标,单位为毫秒:

SparkMetrics_CL
| where Name_s == "streaming.processingDelay"
| extend ClusterName = ClusterName_s, ProcessingDelayMs = Value_d
| project TimeGenerated, ClusterName, ProcessingDelayMs
| render timechart with (title="消息处理延迟趋势(单位:毫秒)")

4. 输入与输出消息速率对比

按1分钟粒度聚合后直观展示速率匹配度:

SparkMetrics_CL
| where Name_s in ("streaming.inputRecords.perSecond", "streaming.outputRecords.perSecond")
| extend MetricType = case(Name_s == "streaming.inputRecords.perSecond", "输入速率", "输出速率")
| summarize AvgRate = avg(Value_d) by bin(TimeGenerated, 1m), MetricType, ClusterName = ClusterName_s
| render barchart with (title="输入输出速率对比(1分钟粒度)", series=MetricType)

5. 每秒消息处理速率

基于批次完成事件计算实际处理速度(排除空闲等待时间,更贴近作业实际处理能力):

SparkListenerEvent_CL
| where EventType_s == "onBatchCompleted"
| extend BatchDurationSec = totimespan(BatchDuration_s)/1s, ProcessedRecords = todlong(NumInputRecords_d)
| extend ProcessRatePerSec = ProcessedRecords / BatchDurationSec
| project TimeGenerated, ClusterName = ClusterName_s, ProcessRatePerSec
| render timechart with (title="单批次每秒消息处理速率")

使用说明

  • 所有查询默认展示所有集群指标,可自行添加 | where ClusterName_s == "目标集群名" 过滤特定集群
  • 时间范围可在 Log Analytics 界面的时间选择器调整,也可在查询开头添加 | where TimeGenerated > ago(24h) 固定范围
  • 若指标字段名和示例不一致,可先执行 SparkMetrics_CL | distinct Name_s 查看环境中实际上报的流处理指标名称

内容的提问来源于stack exchange,提问作者john

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.24 02:45:07