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

Spark 2.4.7 Structured Streaming监控指标及Kafka批次追踪方案咨询

在PySpark 2.4.7中监控结构化流与追踪Kafka批次(无需第三方监控工具)

一、自定义StreamingQueryListener采集核心指标

PySpark 2.4.7虽无Structured Streaming WebUI,但提供StreamingQueryListener接口,可自定义监听器捕获每个批次的进度数据,提取你需要的inputRate、scheduling delay、processingTime、total delay等指标。

实现步骤

  1. 继承StreamingQueryListener,重写三个核心方法:
    • onQueryStarted:记录查询启动事件
    • onQueryProgress:捕获批次进度,提取核心指标
    • onQueryTerminated:记录查询终止事件
  2. 将监听器注册到SparkSession
  3. 将指标输出到日志、本地文件或HBase(与偏移量存储统一)

代码示例

from pyspark.sql.streaming import StreamingQueryListener
import json
from datetime import datetime

class CustomStreamMetricsListener(StreamingQueryListener):
    def onQueryStarted(self, event):
        print(f"[STREAM START] Query ID: {event.id} | Time: {datetime.now()}")

    def onQueryProgress(self, event):
        progress = event.progress
        # 提取核心监控指标
        metrics = {
            "batch_id": progress.batchId,
            "timestamp": datetime.now().isoformat(),
            "input_rows_per_sec": progress.inputRowsPerSecond,  # 对应inputRate
            "processing_time_ms": progress.processingTime,
            "total_delay_ms": progress.totalDelay,
            "scheduling_delay_ms": progress.schedulingDelay  # PySpark 2.4.7原生支持该字段
        }

        # 提取Kafka批次偏移量信息
        kafka_source = progress.sources[0]
        metrics.update({
            "kafka_start_offsets": kafka_source.startOffset,
            "kafka_end_offsets": kafka_source.endOffset
        })

        # 输出到控制台(可替换为写入HBase/本地文件)
        print(f"[STREAM METRICS] {json.dumps(metrics)}")

    def onQueryTerminated(self, event):
        reason = event.exception if event.exception else "正常终止"
        print(f"[STREAM STOP] Query ID: {event.id} | Reason: {reason}")

# 注册监听器
spark.streams.addListener(CustomStreamMetricsListener())

# 你的流处理逻辑示例
kafka_stream = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "broker1:9092,broker2:9092") \
    .option("subscribe", "your_topic") \
    .load()

# 转换与写入SQL Server的逻辑
write_query = kafka_stream.writeStream \
    .format("jdbc") \
    .option("url", "jdbc:sqlserver://sqlserver:1433;databaseName=your_db") \
    .option("dbtable", "target_table") \
    .option("user", "db_user") \
    .option("password", "db_pass") \
    .option("checkpointLocation", "hdfs://path/to/checkpoint") \
    .start()

write_query.awaitTermination()

二、主动查询批次进度

除被动监听,还可通过StreamingQuery对象的lastProgress()方法主动获取最新批次指标,适合定时采样或自定义监控逻辑。

代码示例(后台线程采样)

import threading
import time

def query_progress_monitor(query):
    while query.isActive:
        latest_progress = query.lastProgress()
        if latest_progress:
            print(f"[ACTIVE MONITOR] Batch {latest_progress['batchId']} | Input Rate: {latest_progress['inputRowsPerSecond']} rows/s | Total Delay: {latest_progress['totalDelay']}ms")
        time.sleep(15)  # 每15秒采样一次

# 启动监控线程
monitor_thread = threading.Thread(target=query_progress_monitor, args=(write_query,))
monitor_thread.daemon = True
monitor_thread.start()

三、解析Spark日志提取指标

Spark默认会将结构化流批次进度以INFO级别输出到日志,只需配置日志级别为INFO,即可通过脚本解析获取所有指标。

配置日志级别

修改Spark的log4j.properties,添加:

log4j.logger.org.apache.spark.sql.execution.streaming.StreamExecution=INFO

日志解析脚本(Python)

import re
import json
import csv

log_path = "/var/log/spark/spark-streaming.log"
progress_regex = re.compile(r'INFO  org.apache.spark.sql.execution.streaming.StreamExecution: Query .* made progress: ({.*})')
output_csv = "stream_metrics.csv"

# 初始化CSV文件
with open(output_csv, 'w', newline='') as f:
    writer = csv.writer(f)
    writer.writerow(["batch_id", "timestamp", "input_rate", "processing_time", "total_delay", "scheduling_delay", "kafka_start_offsets", "kafka_end_offsets"])

# 解析日志并写入CSV
with open(log_path, 'r') as f:
    for line in f:
        match = progress_regex.search(line)
        if match:
            progress_data = json.loads(match.group(1))
            batch_id = progress_data["batchId"]
            timestamp = progress_data["timestamp"]
            input_rate = progress_data["inputRowsPerSecond"]
            processing_time = progress_data["processingTime"]
            total_delay = progress_data["totalDelay"]
            scheduling_delay = progress_data["schedulingDelay"]
            kafka_start = json.dumps(progress_data["sources"][0]["startOffset"])
            kafka_end = json.dumps(progress_data["sources"][0]["endOffset"])

            with open(output_csv, 'a', newline='') as f:
                writer = csv.writer(f)
                writer.writerow([batch_id, timestamp, input_rate, processing_time, total_delay, scheduling_delay, kafka_start, kafka_end])

四、Kafka批次追踪

每个批次的batchId与Kafka偏移量范围一一对应,可通过以下方式实现追踪:

  • 在自定义监听器中,将batchId、kafka_start_offsets、kafka_end_offsets与指标一起存储(如写入HBase监控表)
  • 日志解析时,将偏移量信息与批次ID关联保存
  • 排查特定Kafka消息时,可通过偏移量反查对应的batchId,再查看该批次的处理状态与指标

注意事项

  • Kafka偏移量以字典形式存储,包含每个分区的起始/结束偏移,需序列化后存储(如JSON格式)
  • 可复用现有HBase连接,将监控数据与偏移量数据放在同一命名空间,方便关联查询

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 10:37:15