如何将Spark Job类返回值传递至Spark Launcher类
如何从Yarn集群模式的Spark主类向Launcher传递结果?
嘿,这个问题我之前也碰到过——在Yarn集群模式下,Spark应用的Driver是跑在集群节点上的,和你的Launcher进程完全不在同一个JVM(甚至可能不在同一台机器),所以想直接从com.test.Foo返回值给Launcher确实没法像本地方法调用那样直接做,但有几个不用写临时文件的靠谱方案:
方案1:用共享存储传递结果(推荐)
找一个Launcher和Spark Driver都能访问的共享存储(比如Redis、ZooKeeper、关系型数据库都行),在com.test.Foo里把计算结果存入存储,然后Launcher监听Spark应用的状态,一旦应用执行完成,就去共享存储里读取结果。
举个Redis的例子:
在com.test.Foo中的代码
import org.apache.spark.SparkContext; import redis.clients.jedis.Jedis; public class Foo { public static void main(String[] args) { // 执行你的Spark计算逻辑 String calculationResult = "这里是你的计算结果"; // 获取当前Spark应用ID,作为结果的唯一标识 String appId = SparkContext.getOrCreate().applicationId(); // 连接Redis并写入结果 try (Jedis jedis = new Jedis("你的Redis地址", 6379)) { jedis.set("spark_app_result:" + appId, calculationResult); } // 停止Spark上下文 SparkContext.getOrCreate().stop(); } }
在Launcher中的代码
import org.apache.spark.launcher.SparkAppHandle; import org.apache.spark.launcher.SparkLauncher; import redis.clients.jedis.Jedis; public class SparkJobLauncher { public static void main(String[] args) throws Exception { SparkAppHandle handler = new SparkLauncher() .setAppResource("<你的jar包路径>") .setMaster("yarn-cluster") .setDeployMode("cluster") .setVerbose(true) .setMainClass("com.test.Foo") .addAppArgs(args[0], args[1]) .startApplication(new SparkAppHandle.Listener() { @Override public void stateChanged(SparkAppHandle handle) { // 当应用进入最终状态(完成/失败等)时读取结果 if (handle.getState().isFinal()) { String appId = handle.getAppId(); if (appId != null) { try (Jedis jedis = new Jedis("你的Redis地址", 6379)) { String result = jedis.get("spark_app_result:" + appId); if (result != null) { System.out.println("从Spark应用获取到结果:" + result); // 这里可以处理结果 } } } } } @Override public void infoChanged(SparkAppHandle handle) { // 可选:处理应用信息变更(比如日志地址更新) } }); // 保持Launcher进程存活,直到应用完成 while (!handler.getState().isFinal()) { Thread.sleep(1000); } } }
这个方案的优点是可靠性高,适合传递复杂结构的结果,缺点是需要额外部署共享存储服务。
方案2:通过日志解析提取结果
如果不想部署额外服务,可以在com.test.Foo里把结果用唯一标记打印到标准输出或日志中,然后Launcher在应用完成后获取Yarn日志,解析出带标记的结果。
在com.test.Foo中的代码
import org.apache.spark.SparkContext; public class Foo { public static void main(String[] args) { // 计算逻辑 String calculationResult = "你的计算结果"; // 用唯一标记包裹结果,方便后续解析 System.out.println("[SPARK_CUSTOM_RESULT] " + calculationResult); SparkContext.getOrCreate().stop(); } }
在Launcher中的代码
import org.apache.spark.launcher.SparkAppHandle; import org.apache.spark.launcher.SparkLauncher; import java.io.BufferedReader; import java.io.InputStreamReader; public class SparkJobLauncher { public static void main(String[] args) throws Exception { SparkAppHandle handler = new SparkLauncher() .setAppResource("<你的jar包路径>") .setMaster("yarn-cluster") .setDeployMode("cluster") .setVerbose(true) .setMainClass("com.test.Foo") .addAppArgs(args[0], args[1]) .startApplication(new SparkAppHandle.Listener() { @Override public void stateChanged(SparkAppHandle handle) { if (handle.getState() == SparkAppHandle.State.FINISHED) { String appId = handle.getAppId(); if (appId != null) { try { // 执行Yarn日志命令获取应用日志 Process process = new ProcessBuilder("yarn", "logs", "-applicationId", appId) .redirectErrorStream(true) .start(); BufferedReader reader = new BufferedReader(new InputStreamReader(process.getInputStream())); String line; while ((line = reader.readLine()) != null) { // 查找带标记的结果行 if (line.startsWith("[SPARK_CUSTOM_RESULT]")) { String result = line.substring("[SPARK_CUSTOM_RESULT] ".length()); System.out.println("解析到Spark应用结果:" + result); break; } } reader.close(); process.waitFor(); } catch (Exception e) { e.printStackTrace(); } } } } @Override public void infoChanged(SparkAppHandle handle) {} }); while (!handler.getState().isFinal()) { Thread.sleep(1000); } } }
这个方案的优点是无需额外服务,缺点是如果日志量大,解析效率不高,而且如果日志中出现相同标记的内容可能会导致解析错误。
为什么不能直接存入Launcher的当前会话?
要明确的是:Yarn集群模式下,Spark Driver是作为独立的进程在集群节点上启动的,和你的Launcher进程完全隔离——它们不在同一个JVM,内存空间完全不共享。所以你没法把com.test.Foo里的值直接放到Launcher的当前会话中,必须通过跨进程的通信方式(比如上述的共享存储、日志解析)来传递结果。
内容的提问来源于stack exchange,提问作者Tirthankar
相关产品推荐
相关产品推荐

