监控与测量

监控 Spark 应用程序有多种方法:Web UI、指标和外部测量工具。

Web 界面

每个 SparkContext 都会启动一个 Web UI(默认端口为 4040),用于显示有关应用程序的有用信息。其中包括:

您只需在浏览器中打开 http://<driver-node>:4040 即可访问此界面。如果同一主机上运行了多个 SparkContext,它们将绑定到从 4040 开始的后续端口(4041、4042 等)。

请注意,默认情况下,这些信息仅在应用程序运行期间可用。要在事后查看 Web UI,请在启动应用程序之前将 spark.eventLog.enabled 设置为 true。这将配置 Spark 将 UI 中显示的信息记录到持久存储中。

事后查看

只要应用程序的事件日志存在,仍然可以通过 Spark 的历史服务器构建应用程序的 UI。您可以通过执行以下命令启动历史服务器:

./sbin/start-history-server.sh

默认情况下,这会创建一个位于 http://<server-url>:18080 的 Web 界面,列出未完成和已完成的应用程序及尝试。

使用文件系统提供程序类(请参阅下文的 spark.history.provider)时,必须在 spark.history.fs.logDirectory 配置选项中提供基础日志目录,且该目录应包含分别代表每个应用程序事件日志的子目录。

Spark 作业本身必须配置为记录事件,并将其记录到相同的共享可写目录中。例如,如果服务器配置的日志目录为 hdfs://namenode/shared/spark-logs,则客户端选项应为:

spark.eventLog.enabled true
spark.eventLog.dir hdfs://namenode/shared/spark-logs

历史服务器可以按如下方式配置:

环境变量

环境变量含义
SPARK_DAEMON_MEMORY 分配给历史服务器的内存(默认:1g)。
SPARK_DAEMON_JAVA_OPTS 历史服务器的 JVM 选项(默认:无)。
SPARK_DAEMON_CLASSPATH 历史服务器的类路径(默认:无)。
SPARK_PUBLIC_DNS 历史服务器的公共地址。如果未设置,指向应用程序历史记录的链接可能会使用服务器的内部地址,从而导致链接失效(默认:无)。
SPARK_HISTORY_OPTS 历史服务器的 spark.history.* 配置选项(默认:无)。

对滚动事件日志文件应用压缩

长运行应用程序(例如流处理)可能会产生一个巨大的单个事件日志文件,这可能导致高昂的维护成本,并且在 Spark 历史服务器中每次更新时都需要大量资源进行重放。

启用 spark.eventLog.rolling.enabledspark.eventLog.rolling.maxFileSize 可以让您拥有滚动事件日志文件,而不是单个巨大的事件日志文件,这在某些场景下会有帮助,但仍无法帮助您减少日志的总大小。

Spark 历史服务器可以通过在 Spark 历史服务器上设置配置 spark.history.fs.eventLog.rolling.maxFilesToRetain,对滚动事件日志文件进行压缩,从而减小日志的总大小。

具体细节将在下文描述,但请提前注意,压缩是**有损**操作。压缩将丢弃一些在 UI 上不再可见的事件——在启用该选项之前,您可能需要检查哪些事件会被丢弃。

当执行压缩时,历史服务器会列出该应用程序的所有可用事件日志文件,并将索引小于要保留的最小索引文件的事件日志文件视为压缩目标。例如,如果应用程序 A 有 5 个事件日志文件,且 spark.history.fs.eventLog.rolling.maxFilesToRetain 设置为 2,则前 3 个日志文件将被选中进行压缩。

一旦选中目标,它会进行分析以确定哪些事件可以被排除,并将其重写为一个紧凑的文件,同时丢弃决定排除的事件。

压缩试图排除指向过时数据的事件。目前,以下描述了候选的可排除事件:

重写完成后,原始日志文件将以尽力而为的方式被删除。历史服务器可能无法删除原始日志文件,但这不会影响历史服务器的操作。

请注意,如果 Spark 历史服务器发现压缩后空间减少不多,则可能不会压缩旧的事件日志文件。对于流式查询,我们通常预期压缩会运行,因为每个微批处理都会触发一个或多个很快完成的作业,但在批处理查询的许多情况下,压缩不会运行。

另请注意,这是 Spark 3.0 引入的一项新功能,可能尚未完全稳定。在某些情况下,压缩可能会排除超出您预期的事件,从而导致历史服务器上该应用程序出现一些 UI 问题。请谨慎使用。

Spark 历史服务器配置选项

有关 Spark 历史服务器的安全选项,请参阅 安全 页面。

属性名称 默认值 含义 起始版本
spark.history.provider org.apache.spark.deploy.history.FsHistoryProvider 实现应用程序历史记录后端类的名称。目前只有一个由 Spark 提供的实现,它用于查找存储在文件系统中的应用程序日志。 1.1.0
spark.history.fs.logDirectory file:/tmp/spark-events 对于文件系统历史记录提供程序,指定包含要加载的应用程序事件日志的目录 URL。这可以是本地 file:// 路径、HDFS 路径 hdfs://namenode/shared/spark-logs 或 Hadoop API 支持的其他文件系统的路径。 1.1.0
spark.history.fs.update.interval 10s 文件系统历史记录提供程序在日志目录中检查新日志或更新日志的时间间隔。较短的时间间隔可以更快地检测到新应用程序,但代价是服务器在重新读取已更新应用程序时负载更高。更新完成后,已完成和未完成应用程序的列表将立即反映这些变化。 1.4.0
spark.history.retainedApplications 50 缓存中保留 UI 数据的应用程序数量。如果超过此上限,最旧的应用程序将从缓存中移除。如果应用程序不在缓存中,则在从 UI 访问时必须从磁盘加载。 1.0.0
spark.history.ui.maxApplications Int.MaxValue 历史摘要页面上显示的应用程序数量。即使应用程序未显示在历史摘要页面上,也可以通过直接访问其 URL 来使用应用程序 UI。 2.0.1
spark.history.ui.port 18080 历史服务器 Web 界面绑定的端口。 1.0.0
spark.history.kerberos.enabled false 指示历史服务器是否应使用 Kerberos 进行登录。如果历史服务器正在访问安全 Hadoop 集群上的 HDFS 文件,则需要此项。 1.0.1
spark.history.kerberos.principal (无) spark.history.kerberos.enabled=true 时,指定历史服务器的 Kerberos 主体名称。 1.0.1
spark.history.kerberos.keytab (无) spark.history.kerberos.enabled=true 时,指定历史服务器的 Kerberos keytab 文件位置。 1.0.1
spark.history.fs.cleaner.enabled false 指定历史服务器是否应定期从存储中清理事件日志。 1.4.0
spark.history.fs.cleaner.interval 1d spark.history.fs.cleaner.enabled=true 时,指定文件系统作业历史记录清理程序检查要删除的文件的时间频率。如果满足以下两个条件之一,文件将被删除。首先,如果文件的存留时间超过 spark.history.fs.cleaner.maxAge,它们会被删除。其次,如果文件数量超过 spark.history.fs.cleaner.maxNum,Spark 会尝试根据应用程序的最旧尝试时间顺序清理已完成的尝试。 1.4.0
spark.history.fs.cleaner.maxAge 7d spark.history.fs.cleaner.enabled=true 时,文件系统历史记录清理程序运行时,比此时间更早的作业历史记录文件将被删除。 1.4.0
spark.history.fs.cleaner.maxNum Int.MaxValue spark.history.fs.cleaner.enabled=true 时,指定事件日志目录中的最大文件数。Spark 会尝试清理已完成的尝试日志,以将日志目录保持在此限制之下。这应该小于底层文件系统的限制,例如 HDFS 中的 `dfs.namenode.fs-limits.max-directory-items`。 3.0.0
spark.history.fs.endEventReparseChunkSize 1m 在日志文件末尾解析多少字节以查找结束事件。这用于通过跳过事件日志文件中不必要的部分来加快应用程序列表的生成。可以通过将此配置设置为 0 来禁用它。 2.4.0
spark.history.fs.inProgressOptimization.enabled true 启用对进行中的日志的优化处理。此选项可能会使未能重命名其事件日志的已完成应用程序被列为“进行中”。 2.4.0
spark.history.fs.driverlog.cleaner.enabled spark.history.fs.cleaner.enabled 指定历史服务器是否应定期从存储中清理驱动程序日志。 3.0.0
spark.history.fs.driverlog.cleaner.interval spark.history.fs.cleaner.interval spark.history.fs.driverlog.cleaner.enabled=true 时,指定文件系统驱动程序日志清理程序检查要删除的文件的时间频率。只有在文件比 spark.history.fs.driverlog.cleaner.maxAge 更旧时,它们才会被删除。 3.0.0
spark.history.fs.driverlog.cleaner.maxAge spark.history.fs.cleaner.maxAge spark.history.fs.driverlog.cleaner.enabled=true 时,驱动程序日志清理程序运行时,比此时间更早的驱动程序日志文件将被删除。 3.0.0
spark.history.fs.numReplayThreads 可用核心数的 25% 历史服务器用于处理事件日志的线程数。 2.0.0
spark.history.fs.numCompactThreads 可用核心数的 25% 历史服务器用于压缩事件日志的线程数。 4.1.0
spark.history.store.maxDiskUsage 10g 存储缓存应用程序历史信息的本地目录的最大磁盘使用量。 2.3.0
spark.history.store.path (无) 缓存应用程序历史数据的本地目录。如果设置,历史服务器会将应用程序数据存储在磁盘上,而不是保留在内存中。如果历史服务器重启,写入磁盘的数据将被重复使用。 2.3.0
spark.history.store.serializer JSON 用于将内存中 UI 对象写入磁盘(或从磁盘读取)KV 存储的序列化程序;JSON 或 PROTOBUF。JSON 序列化程序是 Spark 3.4.0 之前的唯一选择,因此它是默认值。与 JSON 序列化程序相比,PROTOBUF 序列化程序速度更快且更紧凑。 3.4.0
spark.history.custom.executor.log.url (无) 指定自定义 Spark 执行器日志 URL,以便在历史服务器中支持外部日志服务,而不是使用集群管理器的应用程序日志 URL。Spark 将通过可能因集群管理器而异的模式支持一些路径变量。请查看您的集群管理器文档,了解支持哪些模式(如果有)。此配置对正在运行的应用程序没有影响,它仅影响历史服务器。

目前,仅 YARN 模式支持此配置。

3.0.0
spark.history.custom.executor.log.url.applyIncompleteApplication true 指定是否也将自定义 Spark 执行器日志 URL 应用于未完成的应用程序。如果应将正在运行的应用程序的执行器日志作为原始日志 URL 提供,请将其设置为 `false`。请注意,未完成的应用程序可能包括未正常关闭的应用程序。即使将其设置为 `true`,此配置对正在运行的应用程序也没有影响,它仅影响历史服务器。 3.0.0
spark.history.fs.eventLog.rolling.maxFilesToRetain Int.MaxValue 将被保留为非压缩状态的最大事件日志文件数。默认情况下,将保留所有事件日志文件。出于技术原因,最低值为 1。
请阅读“对旧事件日志文件应用压缩”一节以获取更多详细信息。
3.0.0
spark.history.fs.eventLog.rolling.onDemandLoadEnabled true 在列出文件之前,是否以按需方式查找滚动事件日志位置。 4.1.0
spark.history.store.hybridStore.enabled false 在解析事件日志时是否使用 HybridStore 作为存储。HybridStore 将首先把数据写入内存存储,并有一个后台线程在内存存储写入完成后将数据转储到磁盘存储。 3.1.0
spark.history.store.hybridStore.maxMemoryUsage 2g 可用于创建 HybridStore 的最大内存空间。HybridStore 共同使用堆内存,因此如果启用了 HybridStore,则应通过 SHS 的内存选项增加堆内存。 3.1.0
spark.history.store.hybridStore.diskBackend ROCKSDB 指定混合存储中使用的基于磁盘的存储;ROCKSDB 或 LEVELDB(已弃用)。 3.3.0
spark.history.fs.update.batchSize Int.MaxValue 指定更新新事件日志文件的批处理大小。这可以控制每个扫描过程在合理时间内完成,从而防止初始扫描运行时间过长,并阻塞对新事件日志文件的及时扫描(在大规模环境中)。 3.4.0

请注意,在所有这些 UI 中,表格都可以通过单击标题进行排序,从而轻松识别慢任务、数据倾斜等。

注意

  1. 历史服务器显示已完成和未完成的 Spark 作业。如果应用程序在失败后进行了多次尝试,则会显示失败的尝试,以及任何正在进行的未完成尝试或最终成功的尝试。

  2. 未完成的应用程序仅断续更新。更新间隔时间由检查更改文件的时间间隔(spark.history.fs.update.interval)定义。在较大的集群上,更新间隔可能会设置得很大。查看正在运行的应用程序的实际方法是访问其自己的 Web UI。

  3. 那些未将自己注册为已完成但退出的应用程序将被列为“未完成”——即使它们不再运行。如果应用程序崩溃,可能会发生这种情况。

  4. 发出 Spark 作业完成信号的一种方法是显式停止 Spark Context(sc.stop()),或者在 Python 中使用 with SparkContext() as sc: 结构来处理 Spark Context 的设置和拆卸。

REST API

除了在 UI 中查看指标外,它们也可以以 JSON 格式提供。这为开发人员创建新的 Spark 可视化和监控工具提供了简便的方法。JSON 数据对于正在运行的应用程序和历史服务器都可用。端点挂载在 /api/v1。例如,对于历史服务器,它们通常可在 http://<server-url>:18080/api/v1 访问,对于正在运行的应用程序,可在 https://:4040/api/v1 访问。

在 API 中,应用程序通过其应用程序 ID [app-id] 来引用。在 YARN 上运行时,每个应用程序可能有多次尝试,但只有集群模式下的应用程序才有尝试 ID,客户端模式下的应用程序没有。YARN 集群模式下的应用程序可以通过其 [attempt-id] 来识别。在下文列出的 API 中,当在 YARN 集群模式下运行时,[app-id] 实际上是 [base-app-id]/[attempt-id],其中 [base-app-id] 是 YARN 应用程序 ID。

端点含义
/applications 所有应用程序的列表。
?status=[completed|running] 仅列出处于所选状态的应用程序。
?minDate=[date] 列出的最早开始日期/时间。
?maxDate=[date] 列出的最晚开始日期/时间。
?minEndDate=[date] 列出的最早结束日期/时间。
?maxEndDate=[date] 列出的最晚结束日期/时间。
?limit=[limit] 限制列出的应用程序数量。
示例
?minDate=2015-02-10
?minDate=2015-02-03T16:42:40.000GMT
?maxDate=2015-02-11T20:41:30.000GMT
?minEndDate=2015-02-12
?minEndDate=2015-02-12T09:15:10.000GMT
?maxEndDate=2015-02-14T16:30:45.000GMT
?limit=10
/applications/[app-id]/jobs 给定应用程序的所有作业列表。
?status=[running|succeeded|failed|unknown] 仅列出处于特定状态的作业。
/applications/[app-id]/jobs/[job-id] 给定作业的详细信息。
/applications/[app-id]/stages 给定应用程序的所有阶段列表。
?status=[active|complete|pending|failed] 仅列出处于给定状态的阶段。
?details=true 列出包含任务数据的所有阶段。
?taskStatus=[RUNNING|SUCCESS|FAILED|KILLED|PENDING] 仅列出具有指定任务状态的任务。查询参数 taskStatus 仅在 details=true 时生效。它也支持多个 taskStatus,例如 ?details=true&taskStatus=SUCCESS&taskStatus=FAILED,这将返回所有匹配指定任务状态之一的任务。
?withSummaries=true 列出具有任务指标分布和执行器指标分布的阶段。
?quantiles=0.0,0.25,0.5,0.75,1.0 使用给定的分位数汇总指标。查询参数 quantiles 仅在 withSummaries=true 时生效。默认值为 0.0,0.25,0.5,0.75,1.0
/applications/[app-id]/stages/[stage-id] 给定阶段的所有尝试列表。
?details=true 列出给定阶段包含任务数据的所有尝试。
?taskStatus=[RUNNING|SUCCESS|FAILED|KILLED|PENDING] 仅列出具有指定任务状态的任务。查询参数 taskStatus 仅在 details=true 时生效。它也支持多个 taskStatus,例如 ?details=true&taskStatus=SUCCESS&taskStatus=FAILED,这将返回所有匹配指定任务状态之一的任务。
?withSummaries=true 列出每次尝试的任务指标分布和执行器指标分布。
?quantiles=0.0,0.25,0.5,0.75,1.0 使用给定的分位数汇总指标。查询参数 quantiles 仅在 withSummaries=true 时生效。默认值为 0.0,0.25,0.5,0.75,1.0
示例
?details=true
?details=true&taskStatus=RUNNING
?withSummaries=true
?details=true&withSummaries=true&quantiles=0.01,0.5,0.99
/applications/[app-id]/stages/[stage-id]/[stage-attempt-id] 给定阶段尝试的详细信息。
?details=true 列出给定阶段尝试的所有任务数据。
?taskStatus=[RUNNING|SUCCESS|FAILED|KILLED|PENDING] 仅列出具有指定任务状态的任务。查询参数 taskStatus 仅在 details=true 时生效。它也支持多个 taskStatus,例如 ?details=true&taskStatus=SUCCESS&taskStatus=FAILED,这将返回所有匹配指定任务状态之一的任务。
?withSummaries=true 列出给定阶段尝试的任务指标分布和执行器指标分布。
?quantiles=0.0,0.25,0.5,0.75,1.0 使用给定的分位数汇总指标。查询参数 quantiles 仅在 withSummaries=true 时生效。默认值为 0.0,0.25,0.5,0.75,1.0
示例
?details=true
?details=true&taskStatus=RUNNING
?withSummaries=true
?details=true&withSummaries=true&quantiles=0.01,0.5,0.99
/applications/[app-id]/stages/[stage-id]/[stage-attempt-id]/taskSummary 给定阶段尝试中所有任务的汇总指标。
?quantiles 使用给定的分位数汇总指标。
示例:?quantiles=0.01,0.5,0.99
/applications/[app-id]/stages/[stage-id]/[stage-attempt-id]/taskList 给定阶段尝试的所有任务列表。
?offset=[offset]&length=[len] 列出给定范围内的任务。
?sortBy=[runtime|-runtime] 对任务进行排序。
?status=[running|success|killed|failed|unknown] 仅列出处于该状态的任务。
示例:?offset=10&length=50&sortBy=runtime&status=running
/applications/[app-id]/executors 给定应用程序的所有活动执行器列表。
/applications/[app-id]/executors/[executor-id]/threads 在给定活动执行器内运行的所有线程的堆栈跟踪。历史服务器不可用。
/applications/[app-id]/allexecutors 给定应用程序的所有(活动和死亡)执行器列表。
/applications/[app-id]/storage/rdd 给定应用程序存储的 RDD 列表。
/applications/[app-id]/storage/rdd/[rdd-id] 给定 RDD 的存储状态详细信息。
/applications/[base-app-id]/logs 将给定应用程序所有尝试的事件日志作为 zip 文件中的文件下载。
/applications/[base-app-id]/[attempt-id]/logs 将特定应用程序尝试的事件日志作为 zip 文件下载。
/applications/[app-id]/streaming/statistics 流上下文的统计信息。
/applications/[app-id]/streaming/receivers 所有流接收器的列表。
/applications/[app-id]/streaming/receivers/[stream-id] 给定接收器的详细信息。
/applications/[app-id]/streaming/batches 所有已保留批次的列表。
/applications/[app-id]/streaming/batches/[batch-id] 给定批次的详细信息。
/applications/[app-id]/streaming/batches/[batch-id]/operations 给定批次的所有输出操作列表。
/applications/[app-id]/streaming/batches/[batch-id]/operations/[outputOp-id] 给定操作和给定批次的详细信息。
/applications/[app-id]/sql 给定应用程序的所有查询列表。
?details=[true (default) | false] 列出/隐藏 Spark 计划节点的详细信息。
?planDescription=[true (default) | false] 当物理计划过大时,按需启用/禁用物理 planDescription
?offset=[offset]&length=[len] 列出给定范围内的查询。
/applications/[app-id]/sql/[execution-id] 给定查询的详细信息。
?details=[true (default) | false] 在给定查询详细信息之外,列出/隐藏指标详细信息。
?planDescription=[true (default) | false] 当物理计划过大时,按需为给定查询启用/禁用物理 planDescription
/applications/[app-id]/environment 给定应用程序的环境详细信息。
/version 获取当前 Spark 版本。

可检索的作业和阶段数量受到独立 Spark UI 相同保留机制的限制;"spark.ui.retainedJobs" 定义了触发作业垃圾回收的阈值,spark.ui.retainedStages 定义了阶段的阈值。请注意,垃圾回收在回放时发生:通过增加这些值并重启历史服务器,可以检索到更多的条目。

执行器任务指标

REST API 以任务执行粒度公开了 Spark 执行器收集的任务指标的值。这些指标可用于性能故障排除和工作负载表征。以下是可用指标及其简短描述的列表:

Spark 执行器任务指标名称 简短描述
executorRunTime 执行器运行此任务所花费的耗时。这包括获取 shuffle 数据的时间。该值以毫秒为单位。
executorCpuTime 执行器运行此任务所花费的 CPU 时间。这包括获取 shuffle 数据的时间。该值以纳秒为单位。
executorDeserializeTime 反序列化此任务所花费的耗时。该值以毫秒为单位。
executorDeserializeCpuTime 执行器反序列化此任务所花费的 CPU 时间。该值以纳秒为单位。
resultSize 此任务作为 TaskResult 传回给驱动程序的字节数。
jvmGCTime JVM 在执行此任务时花费在垃圾回收上的耗时。该值以毫秒为单位。
ConcurrentGCCount 此指标返回已发生的集合总数。它仅适用于 Java 垃圾收集器为 G1 Concurrent GC 的情况。
ConcurrentGCTime 此指标返回以毫秒为单位的累计集合耗时近似值。它仅适用于 Java 垃圾收集器为 G1 Concurrent GC 的情况。
resultSerializationTime 序列化任务结果所花费的耗时。该值以毫秒为单位。
memoryBytesSpilled 此任务溢出到内存中的字节数。
diskBytesSpilled 此任务溢出到磁盘上的字节数。
peakExecutionMemory Shuffle、聚合和连接期间创建的内部数据结构所使用的峰值内存。此累加器的值应大致等于任务中创建的所有此类数据结构的峰值大小之和。对于 SQL 作业,这仅跟踪所有不安全(unsafe)运算符和 ExternalSort。
inputMetrics.* 与从 org.apache.spark.rdd.HadoopRDD 或从持久化数据中读取数据相关的指标。
    .bytesRead 读取的总字节数。
    .recordsRead 读取的总记录数。
outputMetrics.* 与外部写入数据(例如写入分布式文件系统)相关的指标,仅在具有输出的任务中定义。
    .bytesWritten 写入的总字节数
    .recordsWritten 写入的总记录数
shuffleReadMetrics.* 与 shuffle 读取操作相关的指标。
    .recordsRead Shuffle 操作中读取的记录数
    .remoteBlocksFetched Shuffle 操作中获取的远程块数
    .localBlocksFetched Shuffle 操作中获取的本地块数(相对于从远程执行器读取)
    .totalBlocksFetched Shuffle 操作中获取的块数(本地和远程)
    .remoteBytesRead Shuffle 操作中读取的远程字节数
    .localBytesRead Shuffle 操作中从本地磁盘读取的字节数(相对于从远程执行器读取)
    .totalBytesRead Shuffle 操作中读取的字节数(本地和远程)
    .remoteBytesReadToDisk Shuffle 操作中读取到磁盘的远程字节数。大块数据在 shuffle 读取操作中被获取到磁盘,而不是读取到内存(这是默认行为)。
    .fetchWaitTime 任务等待远程 shuffle 块所花费的时间。这仅包括阻塞 shuffle 输入数据的时间。例如,如果块 B 正在被获取,而任务尚未完成块 A 的处理,则不认为其被块 B 阻塞。该值以毫秒为单位。
shuffleWriteMetrics.* 与写入 shuffle 数据操作相关的指标。
    .bytesWritten Shuffle 操作中写入的字节数
    .recordsWritten Shuffle 操作中写入的记录数
    .writeTime 阻塞写入磁盘或缓冲区缓存所花费的时间。该值以纳秒为单位。

执行器指标

执行器级指标作为心跳的一部分从每个执行器发送到驱动程序,以描述执行器本身的性能指标,如 JVM 堆内存、GC 信息。执行器指标值及其每个执行器的测量内存峰值通过 REST API 以 JSON 格式和 Prometheus 格式公开。JSON 端点公开于:/applications/[app-id]/executors,Prometheus 端点公开于:/metrics/executors/prometheus。此外,如果 spark.eventLog.logStageExecutorMetrics 为 true,则执行器内存指标的每阶段聚合峰值会写入事件日志。执行器内存指标也通过基于 Dropwizard 指标库 的 Spark 指标系统公开。以下是可用指标及其简短描述的列表:

执行器级指标名称 简短描述
rddBlocks 此执行器块管理器中的 RDD 块。
memoryUsed 此执行器使用的存储内存。
diskUsed 此执行器用于 RDD 存储的磁盘空间。
totalCores 此执行器中可用的核心数。
maxTasks 此执行器中可同时运行的最大任务数。
activeTasks 当前正在执行的任务数。
failedTasks 此执行器中失败的任务数。
completedTasks 此执行器中完成的任务数。
totalTasks 此执行器中的任务总数(运行中、失败和完成)。
totalDuration JVM 在此执行器中执行任务所花费的耗时。该值以毫秒为单位。
totalGCTime JVM 在此执行器中执行垃圾回收所花费的累计耗时。该值以毫秒为单位。
totalInputBytes 此执行器中的总输入字节数。
totalShuffleRead 此执行器中的总 shuffle 读取字节数。
totalShuffleWrite 此执行器中的总 shuffle 写入字节数。
maxMemory 可用于存储的总内存量(以字节为单位)。
memoryMetrics.* 当前内存指标值
    .usedOnHeapStorageMemory 当前用于存储的堆内内存(以字节为单位)。
    .usedOffHeapStorageMemory 当前用于存储的堆外内存(以字节为单位)。
    .totalOnHeapStorageMemory 用于存储的总可用堆内内存(以字节为单位)。此金额可能会随时间而变化,取决于 MemoryManager 实现。
    .totalOffHeapStorageMemory 用于存储的总可用堆外内存(以字节为单位)。此金额可能会随时间而变化,取决于 MemoryManager 实现。
peakMemoryMetrics.* 内存(和 GC)指标的峰值
    .JVMHeapMemory 用于对象分配的堆内存峰值使用量。堆由一个或多个内存池组成。返回的内存使用量的“已使用”和“已提交”大小是所有堆内存池这些值的总和,而返回的内存使用量的“初始”和“最大”大小代表堆内存的设置,可能不是所有堆内存池的总和。在返回的内存使用量中,已用内存量是既被存活对象又被尚未收集的垃圾对象(如果有)所占用的内存量。
    .JVMOffHeapMemory Java 虚拟机使用的非堆内存的峰值使用量。非堆内存由一个或多个内存池组成。返回的内存使用量的“已使用”和“已提交”大小是所有非堆内存池这些值的总和,而返回的内存使用量的“初始”和“最大”大小代表非堆内存的设置,可能不是所有非堆内存池的总和。
    .OnHeapExecutionMemory 使用的堆内执行内存峰值(以字节为单位)。
    .OffHeapExecutionMemory 使用的堆外执行内存峰值(以字节为单位)。
    .OnHeapStorageMemory 使用的堆内存储内存峰值(以字节为单位)。
    .OffHeapStorageMemory 使用的堆外存储内存峰值(以字节为单位)。
    .OnHeapUnifiedMemory 堆内统一内存峰值(执行和存储)。
    .OffHeapUnifiedMemory 堆外统一内存峰值(执行和存储)。
    .DirectPoolMemory JVM 用于直接缓冲区池的内存峰值(java.lang.management.BufferPoolMXBean)。
    .MappedPoolMemory JVM 用于映射缓冲区池的内存峰值(java.lang.management.BufferPoolMXBean)。
    .ProcessTreeJVMVMemory 虚拟内存大小(以字节为单位)。如果 spark.executor.processTreeMetrics.enabled 为 true,则启用。
    .ProcessTreeJVMRSSMemory 驻留集大小:进程在实际内存中占用的页面数。这仅指计入文本、数据或堆栈空间的页面。这不包括未按需加载或已交换出去的页面。如果 spark.executor.processTreeMetrics.enabled 为 true,则启用。
    .ProcessTreePythonVMemory Python 的虚拟内存大小(以字节为单位)。如果 spark.executor.processTreeMetrics.enabled 为 true,则启用。
    .ProcessTreePythonRSSMemory Python 的驻留集大小。如果 spark.executor.processTreeMetrics.enabled 为 true,则启用。
    .ProcessTreeOtherVMemory 其他进程类型的虚拟内存大小(以字节为单位)。如果 spark.executor.processTreeMetrics.enabled 为 true,则启用。
    .ProcessTreeOtherRSSMemory 其他进程类型的驻留集大小。如果 spark.executor.processTreeMetrics.enabled 为 true,则启用。
    .MinorGCCount 轻量级 GC 总次数。例如,垃圾收集器可以是 Copy、PS Scavenge、ParNew、G1 Young Generation 等之一。
    .MinorGCTime 轻量级 GC 总耗时。该值以毫秒为单位。
    .MajorGCCount 重量级 GC 总次数。例如,垃圾收集器可以是 MarkSweepCompact、PS MarkSweep、ConcurrentMarkSweep、G1 Old Generation 等之一。
    .MajorGCTime 重量级 GC 总耗时。该值以毫秒为单位。

RSS 和 Vmem 的计算基于 proc(5)

API 版本控制策略

这些端点已严格版本化,以便于在其之上开发应用程序。特别是,Spark 保证:

请注意,即使在检查正在运行的应用程序的 UI 时,仍需要 applications/[app-id] 部分,尽管此时只有一个应用程序可用。例如,要查看正在运行应用程序的作业列表,您将访问 https://:4040/api/v1/applications/[app-id]/jobs。这是为了使两种模式下的路径保持一致。

指标

Spark 有一个可配置的指标系统,基于 Dropwizard 指标库。这允许用户将 Spark 指标报告给各种接收器(Sinks),包括 HTTP、JMX 和 CSV 文件。这些指标由嵌入在 Spark 代码库中的源生成。它们为特定的活动和 Spark 组件提供测量工具。指标系统通过 Spark 期望存在于 $SPARK_HOME/conf/metrics.properties 的配置文件进行配置。可以通过 spark.metrics.conf 配置属性 指定自定义文件位置。除了使用配置文件外,还可以使用带前缀 spark.metrics.conf. 的一组配置参数。默认情况下,驱动程序或执行器指标使用的根命名空间是 spark.app.id 的值。然而,用户通常希望能够跨应用程序跟踪驱动程序和执行器的指标,而使用应用程序 ID(即 spark.app.id)很难做到这一点,因为它会随着应用程序的每次调用而改变。对于此类用例,可以使用 spark.metrics.namespace 配置属性为指标报告指定自定义命名空间。例如,如果用户想将指标命名空间设置为应用程序名称,他们可以将 spark.metrics.namespace 属性设置为类似 ${spark.app.name} 的值。然后,该值由 Spark 进行适当扩展,并用作指标系统的根命名空间。非驱动程序和非执行器指标从不以 spark.app.id 为前缀,spark.metrics.namespace 属性对此类指标也没有任何影响。

Spark 的指标被解耦到对应于 Spark 组件的不同实例中。在每个实例内,您可以配置一组报告指标的接收器。目前支持以下实例:

每个实例可以向零个或多个接收器报告。接收器包含在 org.apache.spark.metrics.sink 包中:

Prometheus Servlet 镜像了由 Metrics Servlet 和 REST API 公开的 JSON 数据,但采用了时间序列格式。以下是对应的 Prometheus Servlet 端点:

组件 端口 JSON 端点 Prometheus 端点
Master 8080 /metrics/master/json/ /metrics/master/prometheus/
Master 8080 /metrics/applications/json/ /metrics/applications/prometheus/
Worker 8081 /metrics/json/ /metrics/prometheus/
Driver 4040 /metrics/json/ /metrics/prometheus/
Driver 4040 /api/v1/applications/{id}/executors/ /metrics/executors/prometheus/

Spark 还支持一个 Ganglia 接收器,由于许可限制,它未包含在默认构建中。

要安装 GangliaSink,您需要执行 Spark 的自定义构建。请注意,通过嵌入此库,您将在 Spark 包中包含 LGPL 许可的代码。对于 sbt 用户,请在构建前设置 SPARK_GANGLIA_LGPL 环境变量。对于 Maven 用户,请启用 -Pspark-ganglia-lgpl 配置文件。除了修改集群的 Spark 构建外,用户应用程序还需要链接到 spark-ganglia-lgpl 工件。

指标配置文件的语法以及每个接收器可用的参数定义在示例配置文件 $SPARK_HOME/conf/metrics.properties.template 中。

当使用 Spark 配置参数而不是指标配置文件时,相关的参数名称由前缀 spark.metrics.conf. 后跟配置详细信息组成,即参数采用以下形式:spark.metrics.conf.[instance|*].sink.[sink_name].[parameter_name]。此示例显示了 Graphite 接收器的 Spark 配置参数列表。

"spark.metrics.conf.*.sink.graphite.class"="org.apache.spark.metrics.sink.GraphiteSink"
"spark.metrics.conf.*.sink.graphite.host"="graphiteEndPoint_hostName>"
"spark.metrics.conf.*.sink.graphite.port"=<graphite_listening_port>
"spark.metrics.conf.*.sink.graphite.period"=10
"spark.metrics.conf.*.sink.graphite.unit"=seconds
"spark.metrics.conf.*.sink.graphite.prefix"="optional_prefix"
"spark.metrics.conf.*.sink.graphite.regex"="optional_regex_to_send_matching_metrics"

Spark 指标配置的默认值如下:

"*.sink.servlet.class" = "org.apache.spark.metrics.sink.MetricsServlet"
"*.sink.servlet.path" = "/metrics/json"
"master.sink.servlet.path" = "/metrics/master/json"
"applications.sink.servlet.path" = "/metrics/applications/json"

可以使用指标配置文件或配置参数 spark.metrics.conf.[component_name].source.jvm.class=[source_name] 来配置额外的源。目前,JVM 源是唯一可用的可选源。例如,以下配置参数激活了 JVM 源:"spark.metrics.conf.*.source.jvm.class"="org.apache.spark.metrics.source.JvmSource"

可用指标提供程序列表

Spark 使用的指标有多种类型:gauge(仪表)、counter(计数器)、histogram(直方图)、meter(度量计)和 timer(计时器),请参阅 Dropwizard 库文档以获取详细信息。以下组件和指标列表报告了可用指标的名称及一些详细信息,按组件实例和源命名空间分组。Spark 测量中最常使用的指标类型是仪表和计数器。计数器可以识别,因为它们带有 .count 后缀。计时器、度量计和直方图在列表中有标注,列表的其余元素是仪表类型的指标。绝大多数指标在配置其父组件实例后即可激活,有些指标还需要通过额外的配置参数启用,详细信息在列表中报告。

组件实例 = Driver

这是测量指标最多的组件:

组件实例 = Executor

这些指标由 Spark 执行器公开。

源 = JVM Source

备注

组件实例 = applicationMaster

注:在 YARN 上运行时适用

组件实例 = master

注:在 Spark 独立模式作为 master 运行时适用

组件实例 = ApplicationSource

注:在 Spark 独立模式作为 master 运行时适用

组件实例 = worker

注:在 Spark 独立模式作为 worker 运行时适用

组件实例 = shuffleService

注:适用于 shuffle 服务

高级测量

可以使用多种外部工具来帮助分析 Spark 作业的性能:

Spark 还提供了一个插件 API,以便将自定义测量代码添加到 Spark 应用程序中。有两个配置键可用于将插件加载到 Spark 中:

两者都接受以逗号分隔的类名列表,这些类实现 org.apache.spark.api.plugin.SparkPlugin 接口。存在这两个名称是为了使一个列表可以放置在 Spark 默认配置文件中,从而允许用户轻松地从命令行添加其他插件,而无需覆盖配置文件的列表。重复的插件会被忽略。