【发布时间】: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