如何解决无法在JAAS配置中找到“KafkaClient”条目,系统属性‘java.security.auth.login.config’未设置的问题?

内容来源于 Stack Overflow,并遵循CC BY-SA 3.0许可协议进行翻译与使用

  • 回答 (1)
  • 关注 (0)
  • 查看 (3096)

我正在尝试通过spark结构化流媒体连接到Kafka

spark-shell --master local[1] \
       --files /mypath/jaas_mh.conf \
       --packages org.apache.spark:spark-sql-kafka-0-10_2.11:2.3.0 \
       --conf "spark.driver.extraJavaOptions=-Djava.security.auth.login.config=jaas_mh.conf" \
       --conf "spark.executor.extraJavaOptions=-Djava.security.auth.login.config=jaas_mh.conf" \
       --num-executors 1  --executor-cores 1 

当我尝试以编程方式执行相同操作时

object SparkHelper {
  def getAndConfigureSparkSession() = {
    val conf = new SparkConf()
      .setAppName("Structured Streaming from Message Hub to Cassandra")
      .setMaster("local[1]")
      .set("spark.driver.extraJavaOptions", "-Djava.security.auth.login.config=jaas_mh.conf")
      .set("spark.executor.extraJavaOptions", "-Djava.security.auth.login.config=jaas_mh.conf")

    val sc = new SparkContext(conf)
    sc.setLogLevel("WARN")

    getSparkSession()
  }

  def getSparkSession() : SparkSession = {
    val spark = SparkSession
      .builder()
      .getOrCreate()

    spark.sparkContext.addFile("/mypath/jaas_mh.conf")

    return spark
  }
}

我收到错误:

 Could not find a 'KafkaClient' entry in the JAAS configuration. 
    System property 'java.security.auth.login.config' is not set

该如何解决呢?

提问于
用户回答回答于

试试以下代码:

import org.apache.spark.SparkConf
import org.apache.spark.sql.SparkSession

object Driver extends App {

  val confPath: String = "/Users/arcizon/IdeaProjects/spark/src/main/resources/jaas_mh.conf"

  def getAndConfigureSparkSession(): SparkSession = {
    val conf = new SparkConf()
      .setAppName("Structured Streaming from Message Hub to Cassandra")
      .setMaster("local[1]")
      .set("spark.driver.extraJavaOptions", s"-Djava.security.auth.login.config=$confPath")
      .set("spark.executor.extraJavaOptions", s"-Djava.security.auth.login.config=$confPath")

    getSparkSession(conf)
  }

  def getSparkSession(conf: SparkConf): SparkSession = {
    val spark = SparkSession
      .builder()
      .config(conf)
      .getOrCreate()

    spark.sparkContext.addFile(confPath)

    spark.sparkContext.setLogLevel("WARN")

    spark
  }

  val sparkSession: SparkSession = getAndConfigureSparkSession()

  println(sparkSession.conf.get("spark.driver.extraJavaOptions"))
  println(sparkSession.conf.get("spark.executor.extraJavaOptions"))

  sparkSession.stop()
}

扫码关注云+社区

领取腾讯云代金券

年度创作总结 领取年终奖励