如何通过Spark Session获取其创建类的名称?
解决方案
方案1:利用SparkContext Local Properties + 切面织入(无需修改调用方代码)
SparkContext提供的localProperties()支持存储线程绑定的键值对,结合切面编程(如AspectJ)可拦截SparkSession创建调用,自动注入调用类标识:
- 编写切面拦截SparkSession构建方法:
@Aspect public class SparkSessionCreationAspect { @Around("execution(* org.apache.spark.sql.SparkSession$Builder.getOrCreate(..)) || execution(* org.apache.spark.sql.SparkSession$Builder.build(..))") public Object injectCallerClass(ProceedingJoinPoint joinPoint) throws Throwable { // 从栈轨迹中过滤出业务调用类 StackTraceElement[] stackTrace = Thread.currentThread().getStackTrace(); String callerClassName = null; for (StackTraceElement element : stackTrace) { if (!element.getClassName().startsWith("org.apache.spark.") && !element.getClassName().equals(this.getClass().getName())) { callerClassName = element.getClassName(); break; } } if (callerClassName != null) { Object result = joinPoint.proceed(); if (result instanceof SparkSession) { SparkSession sparkSession = (SparkSession) result; sparkSession.sparkContext().localProperties().setProperty("caller_class", callerClassName); } return result; } return joinPoint.proceed(); } }
- 在QueryListener中读取标识:
public class QueryListener implements QueryExecutionListener { @Override public void onSuccess(String funcName, QueryExecution qe, long durationNs) { SparkSession sparkSession = qe.sparkSession(); String callerClass = sparkSession.sparkContext().localProperties().get("caller_class"); // 使用callerClass作为模块/类标识 } @Override public void onFailure(String funcName, QueryExecution qe, Exception exception) { // 失败场景同样可读取该属性 } }
只需通过JVM参数(如-javaagent)加载AspectJ代理即可生效,完全无需修改业务代码。
方案2:自定义SparkSession工具类(需引导调用方切换入口)
若无法使用切面,可封装工具类统一管理SparkSession创建,自动注入调用类标识:
public class SparkSessionUtils { public static SparkSession getSparkSession() { // 获取调用类全限定名 StackTraceElement[] stackTrace = Thread.currentThread().getStackTrace(); String callerClassName = stackTrace[2].getClassName(); SparkSession sparkSession = SparkSession.builder().getOrCreate(); sparkSession.sparkContext().localProperties().setProperty("caller_class", callerClassName); return sparkSession; } }
业务代码仅需将SparkSession.builder().getOrCreate()替换为SparkSessionUtils.getSparkSession(),无需修改其他逻辑,监听器即可读取到调用类名称。
方案3:栈轨迹直接提取(应急调试方案,不推荐生产)
若上述方案均无法实施,可在QueryListener中直接解析当前线程栈轨迹,过滤业务类作为标识:
public class QueryListener implements QueryExecutionListener { @Override public void onSuccess(String funcName, QueryExecution qe, long durationNs) { StackTraceElement[] stackTrace = Thread.currentThread().getStackTrace(); String callerClass = null; // 按业务包名过滤,调整为实际包前缀 for (StackTraceElement element : stackTrace) { if (element.getClassName().startsWith("com.yourcompany.")) { callerClass = element.getClassName(); break; } } // 使用callerClass作为标识 } }
该方案可靠性低(Spark执行线程可能与Session创建线程不一致)、性能差,仅适合临时调试场景。
内容的提问来源于stack exchange,提问作者Anshul
相关产品推荐
相关产品推荐

