You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过Spark Session获取其创建类的名称?

解决方案

方案1:利用SparkContext Local Properties + 切面织入(无需修改调用方代码)

SparkContext提供的localProperties()支持存储线程绑定的键值对,结合切面编程(如AspectJ)可拦截SparkSession创建调用,自动注入调用类标识:

  1. 编写切面拦截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();
    }
}
  1. 在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.17 20:25:40