Не удалось найти запись «KafkaClient» в конфигурации JAAS. Системное свойство 'java.security.auth.login.config' не установлено - PullRequest
0 голосов
/ 25 августа 2018

Я пытаюсь подключиться к 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

Есть указатели?

1 Ответ

0 голосов
/ 25 августа 2018

Даже в conf вы должны указать полный или относительный путь для файла .conf. Также, когда вы создаете SparkConf, я вижу, что вы не применяете его к текущей SparkSession.

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()
}
Добро пожаловать на сайт PullRequest, где вы можете задавать вопросы и получать ответы от других членов сообщества.
...