Azure Databricks数据集读写日志记录方案咨询
解决方案:追踪DataFrame读写存储的大小指标
一、Spark内部日志的可用信息
- Spark的**事件日志(Event Log)**包含作业读写数据的关键任务级指标,Azure Databricks默认开启该功能。你可以通过自定义Spark监听器捕获这些数据,无需开发人员修改业务代码。
- 具体来说,
SparkListenerTaskEnd事件中,taskMetrics.inputMetrics.bytesRead记录了任务读取的字节数,taskMetrics.outputMetrics.bytesWritten记录了任务写入的字节数,汇总这些数据就能得到整个DataFrame的读写大小。 - 这些指标可以关联到作业ID、任务ID等维度,方便后续分析和存储。
二、Ganglia的局限性
- Ganglia仅监控集群节点的系统级指标(如CPU、内存、磁盘IO、网络流量),无法直接提供Spark作业中DataFrame读写的业务级存储大小数据。它只能反映节点整体的IO负载,无法关联到具体的DataFrame或任务。
三、无需封装器的落地方案
1. 自定义Spark全局监听器
编写继承SparkListener的监听器类,在集群初始化时注册,自动收集所有作业的读写指标:
import org.apache.spark.scheduler._ class StorageMetricsCollector extends SparkListener { override def onTaskEnd(taskEnd: SparkListenerTaskEnd): Unit = { // 提取任务读写字节数 val readBytes = taskEnd.taskMetrics.inputMetrics.map(_.bytesRead).getOrElse(0L) val writeBytes = taskEnd.taskMetrics.outputMetrics.map(_.bytesWritten).getOrElse(0L) // 将数据写入外部存储(如ADLS)或Databricks指标系统 // 可结合作业ID、阶段ID、任务ID做维度区分 val metricsRecord = s"${taskEnd.stageId},${taskEnd.taskInfo.taskId},$readBytes,$writeBytes" // 示例:写入ADLS文件 // dbutils.fs.append("dbfs:/mnt/metrics/storage_metrics.log", metricsRecord + "\n") } } // 集群启动时注册监听器(可通过集群初始化脚本或全局配置笔记本执行) spark.sparkContext.addSparkListener(new StorageMetricsCollector)
该方式对开发人员完全透明,所有作业的读写数据都会被自动捕获。
2. 利用Databricks作业API提取指标
通过Databricks Jobs API拉取作业运行的详细指标,其中包含每个阶段的输入输出字节数。可以定时调用API将数据同步到你的指标存储系统,实现统一监控。
3. 解析Spark UI REST API数据
Spark UI的/api/v1/applications/<app-id>/stages接口返回所有阶段的详细指标,包括inputBytes和outputBytes,可以通过脚本定期解析这些数据并存储。
注意事项
- 针对Kafka摄入场景,
inputMetrics.bytesRead会准确记录从Kafka拉取的字节数,符合你的需求。 - 对于Delta Lake等格式,指标反映的是实际IO字节数(受数据跳过、缓存等优化影响),而非逻辑数据大小。
内容的提问来源于stack exchange,提问作者tablet-sediment-mortar
相关产品推荐
相关产品推荐

