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等指标。
实现步骤
- 继承
StreamingQueryListener,重写三个核心方法:onQueryStarted:记录查询启动事件onQueryProgress:捕获批次进度,提取核心指标onQueryTerminated:记录查询终止事件
- 将监听器注册到SparkSession
- 将指标输出到日志、本地文件或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
相关产品推荐
相关产品推荐

