
2. 监控什么一套面向Spark实时任务的监控指标体系很多人在搭建监控系统时第一反应就是先把Grafana装好、把Prometheus配上然后看CPU、看内存、看网络。不能说这个思路完全错误但它确实把顺序搞反了。监控系统服务的是“任务”不是“机器”。你盯着主机指标看半天可能连作业为什么卡住都找不到原因。真正应该先做的是想清楚Spark作业跑起来之后有哪些环节会出问题出问题时会表现出什么迹象哪些指标能在故障发生前就发出预警把这两三个问题想明白再去选监控工具、定采集方案才是稳妥的路线。以我自己的实践经验来看Spark实时任务的监控体系至少需要覆盖以下几个层面每一层都有它不可替代的作用。2.1 Executor维度JVM内存与GC先盯这两个这一层是Spark作业最容易出问题也最容易被忽视的地方。先说结论Spark作业的绝大多数性能问题最终都会落在Executor的JVM内存和GC垃圾回收上。背了很多数据、Shuffle拉取失败、Task反复重试、整个作业像老牛拉车一样慢这些表面现象的背后十有八九都是内存管理出了问题。监控Executor内存核心看两块堆内内存On-Heap的使用量和GC的频次与耗时。堆内内存又被Spark内部的存储体系分成了三块RDD缓存区Storage、Shuffle和算子执行区Execution、以及预留的User Memory。这三个区域的比例关系由spark.memory.fraction和spark.memory.storageFraction两个配置项共同决定。默认情况下spark.memory.fraction为0.6意味着只有60%的堆内空间留给Spark自己管理剩下的40%留给用户代码里的对象和其他系统开销在这60%里Storage和Execution又是动态博弈的关系默认storageFraction为0.5也就是两者各占30%但可以互相抢占。理解了这套机制你再看监控指标就不会一头雾水了。那具体怎么盯从我的实际经验看比较有效的做法是给以下三个指标设置告警堆内存使用率Executor的堆内存使用量占比超过85%时就要开始警惕了。这时候GC会明显频繁起来任务可能出现假死状态。需要说明的是这里不是看总内存而是看JVM Heap Used与spark.executor.memory配置值之间的比例关系。GC耗时与频率Full GC次数在短时间内大幅攀升或者单次Full GC时间超过1秒都是危险信号。我们可以通过JMXJava管理扩展方式采集Executor的GC信息再在Prometheus里做成速率类指标比如“每分钟Full GC次数”。内存溢出异常这个直接看日志就行但要在监控系统中加入日志自动化采集和告警的通道。这一层对应的监控工具是Prometheus JMX Exporter。它可以在Executor的JVM进程里嵌入一个HTTP端口定时暴露堆内存使用量、GC耗时、GC次数、活跃线程数等信息Prometheus每隔几秒拉取一次就足够了。我在生产环境里的经验是给每个Executor固定一个范围地段JMX端口然后把可用的端口池驼在监控配置文件里这样集群扩缩容后监控配置也能自动跟上。到这里可能有人会问“为什么不直接用Spark自带的监控页面上面不也能看到GC和内存吗”自带的监控UI确实能看到但它有一个关键缺陷它不是持久化的历史数据会被滚动清理而且多作业并行时查询起来也比较费力。用Prometheus做时序存储的好处是资源记录可以长期保留事后做性能分析和趋势对比时会方便很多。2.2 作业维度Streaming Batch Duration与处理延迟Spark Streaming的实时作业在调度上其实是一种“微批处理”模式。它把源源不断的数据流按固定的时间间隔切成一段一段的“微批”然后交给Spark的批处理引擎去执行。这个时间间隔就是所谓的batchDuration批处理间隔。理解了这个机制你就能明白实时性和性能之间的关系了。Batch Processing Time批处理耗时——也就是Spark处理一个微批数据所花的总耗时——是可以直接从 Streaming UI 或 REST API 里拿到的指标。当处理耗时逐渐逼近甚至超过批处理间隔时说明这个作业的数据处理能力已经跟不上数据产生的速度了。数据会在内存和磁盘里持续积压最终导致端到端的处理延迟不断攀升。在我实操的过程中建议至少监控以下三个作业级指标Processing Time每批数据的实际处理耗时单位毫秒。Scheduling Delay任务提交到实际运行之间等待的时间。这个值如果长期在几百毫秒以上说明调度资源已经紧张可能需要给作业增加资源或者缩减并发任务数。Total Delay从数据被接收直到该批次完成处理的总耗时。这一项最适合用来衡量实时作业的端到端延迟。这一层的监控我推荐直接使用Spark Streaming的Metrics系统想办法把streaming.lastCompletedBatch_{批次ID}_processingTime、streaming.lastCompletedBatch_submissionTime这类指标接入到Prometheus里。接法的原理并不复杂在Spark应用中启动一个可用的Streaming Context之后注册一个自定义的StreamingListener把每次批次结束时的Processing Time和Scheduling Delay转成Counter或Gauge上报。也可以直接用Spark内置的Metrics配置开启*.sink.prometheus让Spark进程自带暴露Prometheus指标的端口这样连代理都不用额外写了。2.3 数据源与写入链路Kafka Lag与下游写入吞吐实时作业的上游数据源绝大多数是Kafka下游则可能是Redis、数据库、HBase或对象存储。监控上下游的核心目的不是为了看“数据是否流动”而是为了确认“数据是否在以正确的速度流动”。先说上游Kafka消费者滞后量Consumer Lag简称Lag是实时作业健康度的风向标。当Spark Streaming作业的消费速度满足不了生产速度时Lag会持续上涨。看到Lag上涨有两个常见的排查方向一是当前作业自身处理能力不足需要扩容或调优二是可能某个环节偶发故障比如网络抖动、外部系统变慢导致一段时间内处理速度下降后续在慢慢追赶。我在生产里一般这样设计Kafka Lag的监控单独部署一个小服务或者用现成的Kafka Lag Exporter定时去查Spark作业消费组在Kafka各个分区上的消费进度与最新的生产进度做差值得到分区的Lag值然后上报到Prometheus。设置告警时要注意一个容易被文案误导的细节Kafka Lag不是恒为零才健康分区的数据是有波动的偶发的一两个几百毫秒延迟属于正常现象真正需要告警的是Lag持续增长或者达到一定阈值。再看下游实时作业写下游的吞吐量也是一个重要指标。通过Spark累加器Accumulator或者下游存储系统自身的监控指标可以统计每分钟写入的记录数或字节数。如果发现写入吞吐明显低于上游消费吞吐说明下游可能出现了瓶颈。比如写入Redis时如果是用批量管道方式写入吞吐会高很多如果是一条一条地写吞吐就可能成为整个链路中最短的板。2.4 资源利用率CPU、内存、磁盘与网络这一层听起来像是“主机监控”好像没什么技术含量但它的价值在于能帮助判断横向扩容的时机。举个例子如果集群节点的CPU长期在80%以上同时Executor内存使用率也已接近上限你再怎么调Spark参数都成效有限原因很简单——机器本身已经顶到极限了这时应该考虑扩容。这里的主机指标可以直接用node_exporter来采集包括CPU使用率、空闲内存、磁盘IO、网络吞吐等。在生产环境里这些数据跟作业维度的实时数据在同一套Prometheus里存储非常便于把“机器变慢”和“作业延迟”关联起来进行根因分析。3. 如何建一套可落地的监控系统搭建实操我以一套实际的监控架构为例把搭建过程讲透。整体上分四步走数据采集 → 指标存储 → 可视化 → 告警。下面每一层我都会给出具体配置和踩坑经验。3.1 第1步数据采集层的搭建数据采集是整个监控系统的地基。采集做不好后面可视化再漂亮、告警再智能也只是空中楼阁。先搭Prometheus本身。Prometheus是一个开源的时序数据库同时也是监控数据的采集器和告警引擎。我习惯用Thanos或者VictoriaMetrics做长期存储但最基础的一套可以直接用单节点Prometheus跑起来数据保留期设为30天。在资源比较紧张的场景下每天的数据量大概几个GB单节点完全够用。核心要采集的对象有以下几类第一类节点指标每个集群节点上部署node_exporter推荐用systemd托管启动监听端口默认为9100。Prometheus的采集配置里加一个node采集任务规则如下scrape_configs: - job_name: node_exporter static_configs: - targets: [192.168.1.101:9100, 192.168.1.102:9100, 192.168.1.103:9100]把集群里所有节点的IP和端口列出来即可。需要提醒的是Prometheus的采集频率默认为15秒一次对于主机指标这个频率还好但对于Kafka Lag这类变化较快的指标建议单独设置更短的采集间隔比如5秒。可以在scrape_configs里对每个job单独指定scrape_interval。第二类JVM指标每个Spark Executor进程上挂JMX Exporter。这里的关键一步是启动Spark时通过spark.executor.extraJavaOptions参数传入JMX相关配置比如spark-submit \ --conf spark.executor.extraJavaOptions-javaagent:/path/to/jmx_prometheus_javaagent.jar9091:/path/to/jmx_config.yaml \ --conf spark.driver.extraJavaOptions-javaagent:/path/to/jmx_prometheus_javaagent.jar9092:/path/to/jmx_config.yaml \ --class com.example.RealTimeProcessor \ myapp.jar需要注意几个细节每个Executor进程都要自带JMX Exporter端口需要固定或通过环境变量动态分配。固定端口的好处是Prometheus抓取地址稳定但端口数量要按集群最大Executor数量提前预留否则不够用。动态分配的话虽然灵活但每次节点重启Executor后端口会变Prometheus抓取目标会随之变化。我的建议是如果能管控好端口资源优先使用固定端口范围。JMX配置文件中要开启对java.lang:typeMemory、java.lang:typeGarbageCollector等对象的采集。第三类Spark作业指标Spark应用本身需要通过Metrics系统暴露指标。推荐使用内置的prometheussink。在Spark的conf/metrics.properties里加上*.sink.prometheus.classorg.apache.spark.metrics.sink.PrometheusServlet *.sink.prometheus.path/metrics/prometheus随后在提交作业时通过--conf spark.metrics.conf/path/to/metrics.properties指定配置文件Spark就会在Driver和Executor上各自启动一个HTTP接口暴露指标。Prometheus里配置抓取任务时同时抓Driver一个和Executor多个的/metrics/prometheus接口即可。到此Prometheus的scrape_configs应该能自动发现集群中所有需要采集的指标。不过千万注意一点JMX Exporter的端口和Spark Metrics端口要用不同端口否则会冲突导致采集不上。3.2 第2步Grafana可视化面板有了Prometheus下一步就是搭Grafana来做可视化。在Grafana里创建DashBoard并不复杂核心是写PromQL表达式。我把自己长期在用的几个关键面板贴在下面可以直接套用Executor内存使用率面板sum(jvm_memory_bytes_used{instance~.*executor.*, areaheap}) / sum(jvm_memory_bytes_max{instance~.*executor.*, areaheap})这里的过滤条件是假设在Executor的指标标签里带有executor相关的instance标识。实际环境中需要根据自己在Spark中设置的metric标签调整过滤词。GC次数面板sum(increase(jvm_gc_collection_seconds_count{jobspark_executor}[5m])) by (gc_name)Streaming晚点面板spark_streaming_lastCompletedBatch_processingTimeGrafana的UI虽然可以拖拽但在生产环境里我更推荐把所有面板配置保存为JSON文件用Git做版本管理。这样每次调整都有迹可循团队协作时也不容易出现“谁改坏了面板但找不到记录”的尴尬局面。3.3 第3步Alertmanager告警设置告警是整个监控系统的临门一脚。没有告警的监控系统本质上只是一个“事后数据记录器”无法真正帮我第一时间发现作业问题。Alertmanager是Prometheus生态的标准告警组件它负责接收Prometheus的告警规则推送然后通过邮件、钉钉、Webhook等渠道分发出去。在Prometheus的配置文件里定义告警规则比如groups: - name: spark_alerts rules: - alert: ExecutorHeapMemoryHigh expr: | max_over_time( (jvm_memory_bytes_used{areaheap} / jvm_memory_bytes_max{areaheap})[10m:1m] ) 0.85 for: 5m labels: severity: critical annotations: description: Executor堆内存使用率超过85%持续5分钟 - alert: KafkaConsumerLagHigh expr: | max_over_time(kafka_consumergroup_lag{consumer_groupspark-realtime-job}[10m:1m]) 10000 for: 3m labels: severity: warning annotations: description: Kafka消费延迟超过10000条持续3分钟这里的for: 5m很关键它要求触发的条件持续5分钟才真正产生告警可以有效消除误报。生产环境中告警频次一般控制在每天不超过10条为宜否则大家会慢慢忽视告警信息失去告警的意义。3.4 第4步单机到集群的演进与优化单节点Prometheus能扛得住几十台集群规模的采集但如果集群规模达到百台以上单节点就出现了性能瓶颈。此时就需要把采集和存储分离。演进方案我排过两个方案APrometheus Thanos。Thanos是一个开源系统可以把多个Prometheus实例的数据统一查询并且将历史数据持久化到对象存储如S3。这种方式架构比较清晰查询能力也很强。方案BVictoriaMetrics。它作为Prometheus的替代存储兼容PromQL内置了集群模式对高基数数据支持更好性能也更优。如果追求极致的写入吞吐和压缩率我会选这个方案。如果团队里缺少专门的监控运维人力我更倾向于先用单节点Prometheus撑到集群百台规模再说。到了瓶颈期再平滑迁移这比一开始就上一套重方案靠谱得多。4. 优化实录我在生产环境里的四轮调优前面搭建完监控系统接下来要解决的关键问题就是发现作业异常后怎么优化才能让Spark实时作业跑得更快、更稳这部分我结合生产环境的实际案例按“发现问题→分析问题→动手优化→验证效果”的顺序整理几个最典型也最有价值的场景。4.1 第一轮内存参数调整把GC降下来我先说一个比较常见但又容易踩坑的场景。业务场景是一个Spark Streaming实时消费Kafka数据的作业每5秒一个批次。上线后我们在Grafana上看到Executor的堆内存使用率稳定在95%左右Full GC每两分钟触发一次吞吐量拐头向下任务还有偶发的Executor心跳超时。我盯着JVM GC面板看了很久注意到一个小细节老年代Old Gen内存上升得特别快但只升不降。这说明对象没有被及时回收极可能是有大对象长期存活或者分配的内存比例不对。查看历史数据后又发现作业运行初期堆内存使用率长期处于高位但GC一直不触发说明申请的内存一直没有被合理释放。这个场景典型的优化方向是调整内存比例和并发度。如果业务场景里需要持久化大量RDD或DataFrame应该调高Storage区域比例但如果Scraping之后数据直接交给下游则可以把Storage区域让给Execution。于是我将执行环境参数做了调整--conf spark.memory.fraction0.75 --conf spark.memory.storageFraction0.4spark.memory.fraction提高到0.75后Spark自己管理的内存空间扩大了留给用户代码里小对象的内存变少GC的压力反而更小。storageFraction调低成0.4让Execution区域占据更多内存Shuffle和聚合操作的算子执行效率提升。执行完成后在Grafana看的效果非常直接堆内存使用率从稳定95%降到70%~80%Full GC从两分钟一次降到每15分钟一次单个批次的处理耗时也下降了25%左右。这里领悟的经验是调内存参数时必须知道业务在内存里放的是什么。如果没有仔细分析就盲目把memory.fraction调成0.9一旦某个作业的Shuffle操作非常吃内存反而可能触发Executor频繁OOM得不偿失。4.2 第二轮文件小文件问题治理减少Shuffle压力第二个案例来自一个实时写入HDFS的作业。测试期间一切正常但一旦压测数据量上来Shuffle 阶段就会出现大量的Spill到磁盘的操作Batch Processing Time大幅波动从几百毫秒直接飙到3秒以上。我通过Spark UI的“Shuffle Read”标签页定位到问题某个Shuffle Read的Record大小异常大磁盘溢写量比内存中还多。产生这个问题的典型原因是上游生成了大量的小文件而下游读取时每个小文件都对应一次TaskTask数量过多带来的调度开销和Shuffle成本成了一个瓶颈。治理小文件不能只靠“写完就合并”要在写入前和Shuffle过程中双管齐下。我采用的实操方案是在写入HDFS前先做一次Repartition可以在SQL层通过coalesce或repartition调整分区数保证写入的文件数控制在合理范围一般每个文件128MB左右比较合适。如果使用Spark Structured Streaming可以打开File Sink的自动合并功能通过spark.sql.streaming.schemaInference和spark.sql.streaming.fileSink.log.compact这类配置项来实现。在Shuffle读阶段合理设置spark.sql.adaptive.shuffle.targetPostShuffleInputSize来让Spark动态合并小分区。具体在某参数上我把spark.sql.adaptive.coalescePartitions.parallelismFirst设为false强制AQE在动态合并分区时优先保证每个分区足够大而不是坚持维持原有的并行度。跑完后再看Shuffle溢写量减少了90%Batch Processing Time回到了稳定的700ms左右。4.3 第三轮Kafka拉取性能优化第三个案例是关于上游消费的现象是Kafka Lag居高不下消费组每分钟的消费速率跟生产速率差距拉大。正常排查顺序是先看Executor数理够不够再看并行度够不够。如果并行度够却还是来不及消费就可能是单个Task的处理效率太低而不是并发不足。结合Spark Streaming UI里的Input Rate和Processing Time数据我定位到了几个优化点增大Kafka分区的并行消费度将每个Kafka分区对应一个Executor核心保证拉取性能和Task计算能力匹配。设置spark.streaming.kafka.maxRatePerPartition为合理值让每个分区每秒拉取数据量匹配单个Task的处理能力避免拉得太多导致内存压力过大。调优spark.streaming.backpressure.enabled打开背压机制让Spark根据当前处理能力自动调整拉取速度防止数据猛灌进来之后大批量失败。减少Shuffle次数实时作业里尽量避免不必要的宽依赖。能用map就不用reduceByKey能提前filter的就提前filter。经过调整Kafka Lag从持续高位回落到个位数级别整体端到端延迟从分钟级下降到了秒级。4.4 第四轮网络与磁盘IO瓶颈定位最后一轮调优来自一个“玄学”场景数据量和代码没变但作业突然变慢也没有报错。用Grafana查看指标时发现网络接收速率没有明显下降但磁盘IO的utilization达到了85%以上。经过排查发现是同一批节点上还跑了一个批量数据处理作业它的磁盘读写模式是大量随机IO正好和我实时作业的Shuffle溢写撞在一起导致磁盘成为竞争资源。这种场景下可做的优化有几种方向让Spark作业与批处理作业在物理节点上隔离通过YARN或Kubernetes的节点标签做好调度约束。增大Shuffle溢写时使用的磁盘数量让Spark把溢写数据分散在多块物理盘上降低单盘的压力。设置spark.local.dir时配置多个不同挂载点的磁盘路径逗号分隔即可。降低批处理作业的并发度错峰运行。调整完后磁盘IO util很快回落到60%以下Spark作业的Batch Processing Time也恢复正常。5. 生产级实践经验Alerts与问题排查速查表在监控系统落地并运行了一段时间之后我把生产环境中反复出现的问题和排查步骤整理成了速查表希望能帮同样在做事后涂查的人省走一些弯路。5.1 高频问题速查表现象优先排查项常用命令 / 工具解决方向Executor频繁OOM堆内存比例、实际数据量、GC频率Spark UI、JVM GC日志调大executor内存、调整memory.fractionKafka Lag持续上涨处理耗时、消费者数量、分区数Kafka Lag Exporter、Streaming UI增加并行度、开启背压、优化单个Task逻辑单个Batch处理耗时严重超时Shuffle溢写、Task倾斜Spark UI的SQL/Stage页面增加并行度、动态合并分区、优化大表Join磁盘IO瓶颈同类节点上的并发作业iostat、Grafana磁盘面板增加Shuffle本地盘、节点隔离、错峰调度网络带宽打满Shuffle数据量、上游数据量ifstat、Grafana网络面板开启shuffle压缩、优化序列化器、降低数据传输量作业莫名其妙的延迟变高是不是缓存丢了吗Spark UI的Storage页面检查RDD缓存是否被清除、重新计算缓存策略执行计划里Join特别慢数据倾斜Spark UI的执行计划使用Broadcast Join、调整spark.sql.autoBroadcastJoinThreshold这些都是实际排障时比较高效的切入点。如果你遇到类似问题建议按表中顺序走一遍排查逻辑八成能找到症结。5.2 几条独家经验最后分享几条我踩过多次坑之后才总结出来的经验第一监控系统本身要“先跑起来再完美化”。我最开始搭建监控时想把所有指标都做得很完美结果拖了两个星期系统还没上线。后来切换成“先用最核心的5个指标Executor堆内存、GC、Processing Time、Kafka Lag、节点CPU跑起来然后逐步丰富”的策略监控系统很快就在生产上起到了作用。增量式建设比一次性规划更容易落地。第二告警设置一定要按“持续一段时间才触发”的模型来设计。实时作业的指标天然会有波动如果告警阈值设置得过于敏感很容易频繁收到无意义的告警。告警疲劳会造成团队真正遇到故障时无人在意。我的实践是所有告警都加上for: 5m这样的持续时间条件宁可报警慢一点也不要天天狼来了。第三优化不能只调Spark参数要全链路观察。Spark作业只是大数据链路里的一个环节上游数据源的波动、下游写入系统的性能都会影响作业表现。把目光只锁定在Spark自身参数上往往会白费力气。第四也是比较实用的在Spark作业提交脚本里固定记录每次提交时的参数和对应的监控指标。比如我在spark-submit脚本里加了一段脚本每次提交都会把当前作业的配置、日期写入一个日志文件和Prometheus里的时序数据对应起来。这样在排查性能问题时可以直接对比“参数变更前后的数据表现”大幅提升分析效率。6. 工具箱一套完整的监控技术选型清单如果你打算近期启动Spark实时监控系统建设我列一份技术选型清单附上我自己的评价和使用场景可以参考。用途推荐组件说明与个人心得指标采集Prometheus标准方案生态最完善主机指标采集node_exporter部署简单稳定可靠JVM指标采集JMX Exporter可以嵌入Spark进程成本低Spark作业指标Spark Metrics Prometheus Servlet原生支持无需额外代理时序存储VictoriaMetrics集群版性能好压缩率高适合规模化可视化Grafana无悬念首选社区模板丰富告警Alertmanager与Prometheus原生集成策略灵活Kafka消费延迟监测kafka-lag-exporter 或自研Exporter对Spark Streaming消费组专项检测结合我自己的使用体验在规模不够大的场景下最简约有效的组合就是Prometheus node_exporter JMX Exporter Grafana Alertmanager。这套组合不需要额外引入太重的存储部署简单维护成本低却能覆盖前文说的全部监控维度。如果集群规模扩大或者监控指标基数大幅提升Grafana的查询变慢、Prometheus内存吃紧再考虑引入VictoriaMetrics或Thanos做存储扩展。毕竟“够用”是生产环境的第一原则过早引入重型组件反而会增加维护负担。7. 写在最后的个人体会搭了这么多套Spark实时监控系统我越来越体会到一件事监控系统的价值不在于“看得见”而在于“能定位”和“能反应”。看得见只是第一步关键是在指标异常时能快速定位到问题环节及时采取规避动作。记得刚上手时我曾在凌晨被GC告警折腾得晕头转向盯着Prometheus曲线看了大半天才发现问题是某张维度表频繁全量扫描导致Shuffle量大增和最开始怀疑的参数问题完全不相关。从那以后我给自己定了一个规矩优化之前先把自己变成“数据侦探”——认真看指标曲线把异常时间段和前后的操作日志关联起来再决定动手方向。如果你也是刚起步的Spark从业者建议从今天开始把集群里重要的Spark作业先接入一套最朴素的监控哪怕只是Prometheues Grafana把核心指标记录下来。这比等到出现故障再去临时搭系统要宝贵得多因为很多性能问题的根因分析都需要“历史数据”来佐证——没有历史数据排障只能靠猜而靠猜是走不远的。希望这篇文章能帮你在Spark实时监控的搭建和优化路上少踩几个坑。后续如果你在实际落地过程中遇到其他问题也欢迎随时交流。