【问题标题】:error while loading data to bigquery table from dataproc cluster将数据从 dataproc 集群加载到 bigquery 表时出错
【发布时间】:2021-03-01 22:27:04
【问题描述】:

我有一个在 dataproc 中运行的 spark 作业我想将结果加载到 BigQuery,我知道我必须添加 spark-bigquery 连接器才能将数据保存到 bigquery

  name := "spl_prj"

  version := "0.1"

  scalaVersion := "2.11.12"

  val sparkVersion = "2.3.0"

  conflictManager := ConflictManager.latestRevision

  libraryDependencies ++= Seq(
  "org.apache.spark" %%"spark-core" % sparkVersion % Provided,
  "org.apache.spark" %% "spark-sql" % sparkVersion % Provided ,
  "com.google.cloud.spark" %% "spark-bigquery-with-dependencies" % "0.17.3"
  )

当我构建 jar 并提交作业时,它会出现此错误:

  Exception in thread "main" java.lang.ClassNotFoundException: Failed to find data source: bigquery. Please find packages at http://spark.apache.org/third-party-projects.html
at org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSource(DataSource.scala:639)
at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:190)
at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:164)
at com.renault.datalake.spl_prj.Main$.main(Main.scala:58)
at com.renault.datalake.spl_prj.Main.main(Main.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:890)
at org.apache.spark.deploy.SparkSubmit$.doRunMain$1(SparkSubmit.scala:192)
at org.apache.spark.deploy.SparkSubmit$.submit(SparkSubmit.scala:217)
at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:137)
at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala)

原因:java.lang.ClassNotFoundException: bigquery.DefaultSource

example 中提交作业时我没有添加 jar 的权限 我认为当sbt构建jar时,在编译过程中没有添加连接器,我想运行的snippest code scala spark:

   val spark = SparkSession.builder.config(conf).getOrCreate()
   val bucket = "doc_spk"
   spark.conf.set("temporaryGcsBucket", bucket)
   val sc =spark.sparkContext
   val rddRowString = sc.binaryRecords("gs://bucket/GAR", 120).map(x=>(x.slice(0,17),x.slice(17,20),x.slice(20,120)))
   val df=spark.createDataFrame(rddRowString).toDF("v","data","val_data")
   df.write.format("bigquery")
  .option("table","db.table")
  .save()

【问题讨论】:

  • 你必须建立 fat jar

标签: apache-spark google-bigquery sbt google-cloud-dataproc


【解决方案1】:

使用下面的buil.sbt 文件来构建fat jar 文件。

build.sbt

name := "spl_prj"
version := "0.1"
scalaVersion := "2.11.12"
val sparkVersion = "2.3.0"
conflictManager := ConflictManager.latestRevision

libraryDependencies ++= Seq(
  "org.apache.spark" %%"spark-core" % sparkVersion % Provided,
  "org.apache.spark" %% "spark-sql" % sparkVersion % Provided ,
  "com.google.cloud.spark" %% "spark-bigquery-with-dependencies" % "0.17.3"
)
assemblyMergeStrategy in assembly := {
  case PathList("META-INF","services",xs @ _*) => MergeStrategy.filterDistinctLines
  case PathList("META-INF",xs @ _*) => MergeStrategy.discard
  case _ => MergeStrategy.first
}

assemblyOption in assembly := (assemblyOption in assembly).value.copy(includeScala = false)
assemblyJarName in assembly := s"${name.value}-${version.value}.jar"

创建project/plugins.sbt 文件并添加以下内容。

addSbtPlugin("com.eed3si9n" % "sbt-assembly" % "0.15.0")
addSbtPlugin("com.eed3si9n" % "sbt-buildinfo" % "0.9.0")

运行下面的命令来创建```fat`` jar。

sbt clean compile assembly

注意:您可以根据项目要求调整版本。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-12-14
    • 1970-01-01
    • 1970-01-01
    • 2021-03-10
    • 2016-01-26
    • 1970-01-01
    • 2014-07-27
    • 1970-01-01
    相关资源
    最近更新 更多