如何基于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
相关产品推荐
相关产品推荐

