Spark
Spark 是用于大规模数据处理的分布式计算引擎。采集范围包括 Master 的集群调度状态、Worker 的资源容量、Driver 的调度健康,以及 Executor 的内存、GC 和 Shuffle I/O。
配置¶
前置条件¶
- 已部署 Spark 集群;
- DataKit 能访问各 Spark JVM 暴露的
/metrics地址; - 已下载与 JVM 版本兼容的 JMX Exporter Java Agent;
- Exporter 端口仅在可信网络中开放。
部署 JMX Exporter¶
建议为 Master、Worker、Driver 和 Executor 使用独立的 Exporter 端点和规则文件。规则需匹配目标 Spark MBean,并保留资源归属标签;同一主机运行多个 Worker 或 Executor 时,使用 worker_id、executor_id 区分逻辑实例,而不是只依赖 host。
在 Master 与 Worker 的启动参数中添加 Java Agent。以下使用端口映射示例中的 Master、Worker-1 端口:
# Master
export SPARK_MASTER_OPTS="$SPARK_MASTER_OPTS -javaagent:/opt/jmx/jmx_prometheus_javaagent.jar=28099:/opt/jmx/spark-master.yml"
# Worker
export SPARK_WORKER_OPTS="$SPARK_WORKER_OPTS -javaagent:/opt/jmx/jmx_prometheus_javaagent.jar=28199:/opt/jmx/spark-worker.yml"
提交应用时,可为 Driver 和 Executor 添加 Java Agent:
spark-submit \
--conf 'spark.driver.extraJavaOptions=-javaagent:/opt/jmx/jmx_prometheus_javaagent.jar=24099:/opt/jmx/spark-driver.yml' \
--conf 'spark.executor.extraJavaOptions=-javaagent:/opt/jmx/jmx_prometheus_javaagent.jar=28100:/opt/jmx/spark-executor.yml' \
...
重启对应 Spark 进程后,在 DataKit 所在主机确认端点可访问:
配置 DataKit Prometheus 采集器¶
进入 DataKit 安装目录 下的 conf.d/prom,复制 prom.conf.sample 为 spark.conf。一个配置文件可包含多个 [[inputs.prom]] 采集块,推荐将 Master、Worker、Driver 和 Executor 端点集中管理。
下表为 3 Worker 测试集群的本机端口映射示例。生产环境请替换为实际服务地址或服务发现地址;Executor 的端点数随应用运行状态变化。
| 组件 | Exporter 地址示例 | 采集标签 |
|---|---|---|
| Master | http://127.0.0.1:28099/metrics |
component=master |
| Worker-1 / Worker-2 / Worker-3 | 28199 / 28299 / 28399 端口的 /metrics |
component=worker,并分别设置 worker_id |
| Driver | http://127.0.0.1:24099/metrics |
component=driver |
| Executor-1 / Executor-2 / Executor-3 | 28100 / 28200 / 28300 端口的 /metrics |
component=executor,并设置所属 worker_id |
spark.conf 示例:
# Master
[[inputs.prom]]
urls = ["http://127.0.0.1:28099/metrics"]
source = "spark-master"
measurement_name = "spark"
interval = "10s"
metric_types = []
[inputs.prom.tags]
env = "test"
component = "master"
collector = "jmx-exporter"
# Worker-1;Worker-2、Worker-3 复制此块并修改 url、worker_id
[[inputs.prom]]
urls = ["http://127.0.0.1:28199/metrics"]
source = "spark-worker"
measurement_name = "spark"
interval = "10s"
metric_types = []
[inputs.prom.tags]
env = "test"
component = "worker"
collector = "jmx-exporter"
worker_id = "worker-1"
# Driver
[[inputs.prom]]
urls = ["http://127.0.0.1:24099/metrics"]
source = "spark-driver"
measurement_name = "spark"
interval = "10s"
metric_types = []
[inputs.prom.tags]
env = "test"
component = "driver"
collector = "jmx-exporter"
# Executor-1;Executor-2、Executor-3 复制此块并修改 url、worker_id、executor_id
[[inputs.prom]]
urls = ["http://127.0.0.1:28100/metrics"]
source = "spark-executor"
measurement_name = "spark"
interval = "10s"
metric_types = []
[inputs.prom.tags]
env = "test"
component = "executor"
collector = "jmx-exporter"
worker_id = "worker-1"
executor_id = "1"
urls 指向 JMX Exporter 指标地址;source 用于区分采集器;measurement_name 指定指标写入 spark 指标集。动态创建或退出的 Executor 应通过服务发现或配置管理维护采集块,避免使用同一个静态 executor_id 覆盖多个实例。
可先在 DataKit 所在主机检查端点是否返回预期的 spark_* 指标:
curl -fsS http://127.0.0.1:28099/metrics | rg '^spark_master_'
curl -fsS http://127.0.0.1:24099/metrics | rg '^spark_driver_'
配置完成后,按 DataKit 服务管理 重启 DataKit。
指标¶
Spark 指标默认写入 spark 指标集。实际字段由 JMX Exporter 规则决定,以下为运维中常用字段。
集群与 Worker 容量¶
| 字段 | 含义 | 推荐分组 |
|---|---|---|
master_aliveWorkers |
存活 Worker 数 | host |
master_workers |
已注册 Worker 数 | host |
master_apps |
运行中应用数 | host |
master_waitingApps |
等待调度的应用数 | host |
worker_coresUsed / worker_coresFree |
Worker 已用 / 空闲 Core | worker_id |
worker_memUsed_MB / worker_memFree_MB |
Worker 已用 / 空闲内存 | worker_id |
worker_executors |
Worker 当前运行 Executor 数 | worker_id |
Driver 调度健康¶
| 字段 | 含义 | 单位 |
|---|---|---|
driver_DAGScheduler_job_activeJobs |
活跃 Job 数 | count |
driver_DAGScheduler_stage_runningStages |
运行中 Stage 数 | count |
driver_DAGScheduler_stage_waitingStages |
等待 Stage 数 | count |
driver_DAGScheduler_stage_failedStages |
失败 Stage 数 | count |
driver_LiveListenerBus_queue_*_size |
ListenerBus 队列深度 | count |
driver_LiveListenerBus_queue_*_numDroppedEvents_total |
ListenerBus 累计丢弃事件数 | count |
Executor 资源与 I/O¶
| 字段 | 含义 | 单位 |
|---|---|---|
executor_JVMHeapMemory |
JVM 堆内存 | B |
executor_OnHeapExecutionMemory / executor_OnHeapStorageMemory |
On-Heap 执行 / 存储内存 | B |
executor_threadpool_activeTasks |
活跃任务数 | count |
executor_MajorGCCount / executor_MinorGCCount |
GC 累计次数 | count |
executor_succeededTasks_total |
成功任务累计数 | count |
executor_shuffleTotalBytesRead_total / executor_shuffleBytesWritten_total |
Shuffle 累计读写量 | B |
executor_diskBytesSpilled_total / executor_memoryBytesSpilled_total |
磁盘 / 内存 Spill 累计量 | B |
*_total 表示累计计数器。需展示吞吐或速率时,请先确认计数器持续上报,再用 rate() 计算;累计量、速率、内存和任务数量不应放在同一坐标轴。