Apache Spark是否有捕获SQL源/接收器读写流量的指标?
解决方案:捕获Spark SQL数据库读写流量指标
一、Spark默认指标的局限性
Spark原生的bytesRead和bytesWritten属于Executor Task Metrics,仅针对HDFS、本地文件系统这类基于Hadoop InputFormat/OutputFormat的存储系统统计流量。对于JDBC这类SQL数据库的读写操作,数据传输由数据库驱动完成,Spark默认不会统计这部分字节量,因此没有现成的原生指标可用。
二、自定义实现方案
可以通过自定义JDBC数据源包装结合Spark度量系统/Listener实现SQL读写流量统计,具体步骤如下:
1. 自定义JDBC数据源,统计读写字节数
- 针对JDBC读取:继承
JdbcRelationProvider,构建JdbcRelation时包装ResultSet,遍历结果集时统计每行数据的字节大小(可通过ResultSetMetaData获取字段类型计算字节数,或序列化行数据后统计长度),累计得到总读取字节数。 - 针对JDBC写入:继承
JdbcBatchWrite,执行批量写入时统计每条待写入数据的字节大小,累计得到总写入字节数。
2. 将统计结果注入Spark度量系统
把上述统计的字节数注册为Spark自定义度量:
- 通过
SparkEnv.get.metricsSystem获取度量系统实例,注册Counter类型的度量(比如jdbc_bytes_read、jdbc_bytes_written)。 - 在JDBC读写完成后,更新这些Counter的数值。
3. 适配spark-monitoring采集
由于你使用的spark-monitoring基于log4j2,可通过两种方式让自定义指标被采集:
- 方式一:在自定义数据源/Listener中,将统计的字节量以spark-monitoring可解析的格式输出到日志,让log4j2直接采集并发送到LogAnalytics。
- 方式二:利用Spark的Metrics系统暴露自定义指标,修改spark-monitoring配置,使其采集这些新增指标。
4. 可选:通过SparkListener统一上报
如果需要在任务/阶段结束时统一上报指标,可以实现SparkListener,监听onTaskEnd或onJobEnd事件,从任务的TaskMetrics中提取自定义注入的JDBC流量数据,批量发送到LogAnalytics。
内容的提问来源于stack exchange,提问作者Riccardo
相关产品推荐
相关产品推荐

