Scala项目Spark Submit提交Jar类:复用Context配置与参数传递
解决Scala Spark项目复用spark-submit配置+可选参数的问题
嘿,这个需求很常见,我来给你一步步拆解解决办法:
一、复用spark-submit配置并选择性覆盖参数
其实Spark已经帮你做好了大部分工作——当你用spark-submit提交任务时,所有通过命令行设置的参数(比如--num-executors 6)都会自动加载到Spark的配置里。你不需要手动读取这些参数,直接基于现有配置修改需要覆盖的部分就行。
举个例子,如果你想覆盖应用名称,同时保留spark-submit设置的其他参数:
import org.apache.spark.sql.SparkSession object MySparkApp { def main(args: Array[String]): Unit = { // 基于spark-submit的配置创建SparkSession,同时覆盖应用名称 val spark = SparkSession.builder() .config("spark.app.name", "我的自定义应用名") // 这里会覆盖spark-submit里的--conf spark.app.name .getOrCreate() // 验证一下:打印所有配置,你会看到num-executors是6,应用名称是我们设置的 spark.sparkContext.getConf.getAll.foreach { case (k, v) => println(s"$k = $v") } // 后续业务逻辑... spark.stop() } }
如果想优先保留spark-submit的配置,只有当没设置时才用代码里的默认值,可以这样写:
val sparkConf = new SparkConf() // 仅当spark-submit没设置spark.app.name时,才用默认值 val appName = sparkConf.getOption("spark.app.name").getOrElse("默认应用名") val spark = SparkSession.builder() .config(sparkConf) .config("spark.app.name", appName) .getOrCreate()
二、处理可选命令行参数
推荐用scopt这个轻量级的Scala命令行解析库,它能帮你轻松定义和解析可选参数。
步骤1:添加依赖到build.sbt
libraryDependencies += "com.github.scopt" %% "scopt" % "4.1.0"
步骤2:编写参数解析逻辑
先定义一个case class来存储参数,然后用scopt的OptionParser来处理:
import com.github.scopt.OptionParser // 定义参数结构,默认值对应可选参数的默认情况 case class Config( appName: String = "默认应用名", inputPath: String = "", outputPath: String = "" ) object MySparkApp { def main(args: Array[String]): Unit = { // 创建解析器 val parser = new OptionParser[Config]("MySparkApp") { head("MySparkApp", "1.0") // 应用名称和版本 // 可选参数:--name 或 -n,设置应用名称 opt[String]('n', "name") .action((value, config) => config.copy(appName = value)) .text("可选:自定义应用名称,优先级高于spark-submit的配置") // 必填参数:--input 或 -i,输入路径 opt[String]('i', "input") .required() .action((value, config) => config.copy(inputPath = value)) .text("必填:输入文件路径") // 必填参数:--output 或 -o,输出路径 opt[String]('o', "output") .required() .action((value, config) => config.copy(outputPath = value)) .text("必填:输出文件路径") // 帮助信息 help("help").text("显示帮助信息") } // 解析命令行参数 parser.parse(args, Config()) match { case Some(config) => // 根据参数配置SparkSession val sparkConf = new SparkConf() // 如果用户指定了--name,就用这个值;否则用spark-submit的配置,没有的话用默认 val finalAppName = if (config.appName != "默认应用名") config.appName else sparkConf.getOption("spark.app.name").getOrElse("默认应用名") val spark = SparkSession.builder() .config(sparkConf) .config("spark.app.name", finalAppName) .getOrCreate() // 业务逻辑示例:读取输入,处理,写入输出 val df = spark.read.text(config.inputPath) df.write.text(config.outputPath) spark.stop() case None => // 参数解析失败(比如缺必填参数),scopt会自动打印帮助信息,这里直接退出 sys.exit(1) } } }
三、使用示例
打包成Jar后,这样提交任务:
# 同时设置spark-submit参数和自定义可选参数 spark-submit --num-executors 6 --class org.project.MySparkApp my-spark-app.jar \ --name "我的任务名称" \ --input "/path/to/input.txt" \ --output "/path/to/output"
如果不指定--name,就会用spark-submit里的应用名称(如果有的话),否则用默认值。
内容的提问来源于stack exchange,提问作者ysig
相关产品推荐
相关产品推荐

