【问题标题】:Scala Spark: Multiple sources found for jsonScala Spark:为 json 找到多个来源
【发布时间】:2020-07-05 15:57:14
【问题描述】:

在我的 hadoop 集群上执行 spark2-submit 时遇到异常,当我在 hdfs 中读取 .jsons 的目录时,我不知道如何解决它。

我在几个板上找到了一些关于此的问题,但没有一个很受欢迎或没有答案。

我尝试显式导入 org.apache.spark.sql.execution.datasources.json.JsonFileFormat,但导入 SparkSession 似乎是多余的,因此无法识别。

不过,我可以确认这两个课程都可用。

val json:org.apache.spark.sql.execution.datasources.json.JsonDataSource
val json:org.apache.spark.sql.execution.datasources.json.JsonFileFormat

堆栈跟踪:

Exception in thread "main" org.apache.spark.sql.AnalysisException: Multiple sources found for json (org.apache.spark.sql.execution.datasources.json.JsonFileFormat, org.apache.spark.sql.execution.datasources.json.DefaultSource), please specify the fully qualified class name.;
    at org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSource(DataSource.scala:670)
    at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:190)
    at org.apache.spark.sql.DataFrameReader.json(DataFrameReader.scala:397)
    at org.apache.spark.sql.DataFrameReader.json(DataFrameReader.scala:340)
    at jsonData.HdfsReader$.readJsonToDataFrame(HdfsReader.scala:45)
    at jsonData.HdfsReader$.process(HdfsReader.scala:52)
    at exp03HDFS.StartExperiment03$.main(StartExperiment03.scala:41)
    at exp03HDFS.StartExperiment03.main(StartExperiment03.scala)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:498)
    at org.apache.spark.deploy.JavaMainApplication.start(SparkApplication.scala:52)
    at org.apache.spark.deploy.SparkSubmit$.org$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:894)
    at org.apache.spark.deploy.SparkSubmit$.doRunMain$1(SparkSubmit.scala:198)
    at org.apache.spark.deploy.SparkSubmit$.submit(SparkSubmit.scala:228)
    at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:137)
    at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala)

HdfsReader:

import java.net.URI
import org.apache.hadoop.fs.{LocatedFileStatus, RemoteIterator}
import org.apache.spark.sql.DataFrame
import org.apache.spark.sql.SparkSession
import pipelines.ContentPipeline

object HdfsReader {

... 

  def readJsonToDataFrame(inputDir: String, multiline: Boolean = true, verbose: Boolean = false)
  : DataFrame = {

    val multiline_df = spark.read.option("multiline",value = true).json(inputDir)
    multiline_df.show(false)
    if (verbose) multiline_df.show(truncate = true)
    multiline_df
  }
  
  def process(path: URI) = {
    val dataFrame = readJsonToDataFrame(path.toString, verbose = true)
    val contentDataFrame = ContentPipeline.getContentOfText(dataFrame)
    val newDataFrame = dataFrame.join(contentDataFrame, "text").distinct()
    JsonFileUtils.saveAsJson(newDataFrame, outputFolder)
  }

}

build.sbt

version := "0.1"
scalaVersion := "2.11.8" //same version hadoop uses


libraryDependencies ++=Seq(
  "org.apache.spark" %% "spark-core" % "2.3.0", //same version hadoop uses
  "com.johnsnowlabs.nlp" %% "spark-nlp" % "2.3.0",
  "org.apache.spark" %% "spark-sql" % "2.3.0",
  "org.apache.spark" %% "spark-mllib" % "2.3.0",
  "org.scalactic" %% "scalactic" % "3.2.0",
  "org.scalatest" %% "scalatest" % "3.2.0" % "test",
  "com.lihaoyi" %% "upickle" % "0.7.1")

【问题讨论】:

    标签: apache-spark hadoop apache-spark-sql


    【解决方案1】:

    您的类路径中似乎同时有 Spark 2.x 和 3.x jar。根据 sbt 文件,应该使用 Spark 2.x,但是,JsonFileFormat 在 Spark 3.x 中添加了 this issue

    【讨论】:

    • 我该如何具体解决这个问题?这是否意味着我不能使用 JsonFileFormat?
    • 需要使用 Spark 2 或 3 但不能同时使用,即去掉重复的 jars。也许您的集群安装了 Spark 3?看来spark-nlp 目前只支持 Spark 2
    • 更新了版本以适应集群上的版本
    • 在问题中添加了 sbt evited
    • 我检查了构建 sbt 中的每个版本以与集群兼容,但我仍然得到相同的确切错误
    【解决方案2】:

    所以,我解决了我的问题:

    val dataFrame1 = spark
      .read
      .option("multiLine", value = true)
      .json(inputDir)
    
    val dataFrame2 = spark
      .read
      .format("org.apache.spark.sql.execution.datasources.json.JsonFileFormat")
      .option("multiline",value = true)
      .load(inputDir)
    

    这两个函数的作用基本相同:

    他们将*.json 文件的整个目录读入DataFrame。

    唯一不同的是,dataFrame1 对您要使用的数据类型做出假设,并在 org.apache.spark.sql.execution.datasources.json 中查找它。

    你不希望这样,因为如果你尝试从这个类路径初始化一个 json,你会发现 2 个源。

    val json:org.apache.spark.sql.execution.datasources.json.JsonDataSource
    val json:org.apache.spark.sql.execution.datasources.json.JsonFileFormat
    

    但是有一个option 允许您指定来源,以防发生冲突。

    这是您使用format(source).read(path) 显式使用特定数据类型来读取文件的地方。

    【讨论】:

      猜你喜欢
      • 2021-06-14
      • 2021-04-09
      • 1970-01-01
      • 2021-09-13
      • 1970-01-01
      • 2020-04-24
      • 2021-09-03
      • 2016-07-20
      • 2016-11-21
      相关资源
      最近更新 更多