Spark Streaming能否从Executor直接发送Kafka指标以避免监控遗漏?
解决Spark Streaming无过滤数据时Kafka偏移量指标丢失的问题
针对你遇到的“无符合过滤条件数据时,Spark不发送Prometheus指标,导致无法发现其他生产者过载”的问题,这里提供几个可行的解决思路:
一、读取Kafka后立即上报偏移量指标
不管后续业务过滤是否有数据,在读取Kafka的环节就单独处理偏移量指标上报,不依赖后续的数据库写入动作。
具体操作:
- 读取Kafka时保留偏移量元数据
用spark.readStream读取Kafka时,通过selectExpr把topic、partition、offset这些关键元数据字段一起读出来,示例代码:val kafkaDF = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "你的Kafka broker地址") .option("subscribe", "要订阅的主题") .load() .selectExpr("topic", "partition", "offset", "value") // 保留元数据和原始消息 - 在微批次中提取偏移量并上报
使用foreachBatch在每个微批次里先统计当前批次的最小/最大偏移量,然后推送到Prometheus Pushgateway(因为Spark Executor没法直接被Prometheus拉取指标,Pushgateway是常用的方式)。示例:
注意:要确保Pushgateway在K8S集群内可访问,Spark作业有对应的网络权限,同时需要引入Prometheus Java客户端依赖(提交作业时加import io.prometheus.client.{Gauge, CollectorRegistry} import io.prometheus.client.exporter.PushGateway // 定义Prometheus指标 val kafkaOffsetGauge = Gauge.build() .name("kafka_topic_offset") .labelNames("topic", "partition", "offset_type") // offset_type区分min/max偏移量 .help("Spark Streaming读取的Kafka主题分区偏移量") .register() kafkaDF.writeStream .foreachBatch { (batchDF, batchId) => // 第一步:统计当前批次的偏移量 val offsetStats = batchDF.groupBy("topic", "partition") .agg(min("offset").alias("min_offset"), max("offset").alias("max_offset")) .collect() // 第二步:设置指标值 offsetStats.foreach { row => val topic = row.getAs[String]("topic") val partition = row.getAs[Int]("partition") kafkaOffsetGauge.labels(topic, partition.toString, "min").set(row.getAs[Long]("min_offset")) kafkaOffsetGauge.labels(topic, partition.toString, "max").set(row.getAs[Long]("max_offset")) } // 第三步:推送到Pushgateway val pushGateway = new PushGateway("你的Pushgateway地址:9091") pushGateway.pushAdd(CollectorRegistry.defaultRegistry, "spark-streaming-job") // 最后处理业务逻辑:过滤、写数据库 val filteredDF = batchDF.filter("producer_id = '你的目标生产者ID'") filteredDF.write.format("jdbc").option("url", "数据库地址").option(...).save() } .start() .awaitTermination()--packages io.prometheus:simpleclient_pushgateway:0.16.0)。
二、启用Spark内置的Kafka指标
Spark 3.3.2本身内置了Kafka消费相关的指标,比如kafka_consumer_offsets,只要Spark读取了Kafka批次,不管后续过滤是否有数据,这些指标都会更新。你只需要配置Spark暴露这些指标给Prometheus:
配置步骤:
- 提交作业时添加指标相关配置
spark-submit \ --conf spark.metrics.conf.*.sink.prometheus.class=org.apache.spark.metrics.sink.PrometheusSink \ --conf spark.metrics.conf.*.sink.prometheus.port=9100 \ --conf spark.metrics.conf.*.source.jvm.class=org.apache.spark.metrics.source.JvmSource \ --packages org.apache.spark:spark-metrics-prometheus_2.12:3.3.2 \ --class 你的主类 你的jar包 - 在K8S中配置Prometheus采集
给Spark Driver和Executor的Pod添加标签,然后在Prometheus的采集规则中配置抓取9100端口的指标。之后在Grafana中就能基于kafka_consumer_offsets等指标监控偏移量,不用依赖业务写入动作。
三、用Kafka Exporter独立监控主题偏移量
如果不想修改Spark代码,直接部署Kafka Exporter来监控Kafka集群的所有主题偏移量、消费滞后情况。Kafka Exporter直接从Kafka获取数据,不受Spark作业的影响,能直观看到各个生产者的写入速度、主题的偏移量增长情况,即使Spark没上报指标,也能发现某个生产者过载的问题。
部署要点:
- 使用官方镜像
danielqsj/kafka-exporter:latest在K8S中部署Deployment,配置参数指向你的Kafka broker:args: - --kafka.server=你的Kafka broker地址:9092 - 配置Prometheus抓取Kafka Exporter的默认端口9308,然后在Grafana中导入现成的Kafka Exporter仪表盘模板,就能监控所有主题的偏移量、分区滞后等关键指标。
内容的提问来源于stack exchange,提问作者Александр Трутнев
相关产品推荐
相关产品推荐

