【问题标题】:How to write and update by kudu API in Spark 2.1如何在 Spark 2.1 中通过 kudu API 编写和更新
【发布时间】:2018-01-10 16:28:52
【问题描述】:

我想通过 Kudu API 编写和更新。 这是maven依赖:

<dependency>
  <groupId>org.apache.kudu</groupId>
  <artifactId>kudu-client</artifactId>
  <version>1.1.0</version>
</dependency>
<dependency>
  <groupId>org.apache.kudu</groupId>
  <artifactId>kudu-spark2_2.11</artifactId>
  <version>1.1.0</version>
</dependency>

在下面的代码中,我对KuduContext参数一无所知。

我在 spark2-shell 中的代码:

val kuduContext = new KuduContext("master:7051") 

Spark 2.1 流中也出现同样的错误:

import org.apache.kudu.spark.kudu._
import org.apache.kudu.client._
val sparkConf = new SparkConf().setAppName("DirectKafka").setMaster("local[*]")
val ssc = new StreamingContext(sparkConf, Seconds(2))
val messages = KafkaUtils.createDirectStream("")
messages.foreachRDD(rdd => {
   val spark = SparkSession.builder.config(rdd.sparkContext.getConf).getOrCreate()
   import spark.implicits._
   val bb = spark.read.options(Map("kudu.master" -> "master:7051","kudu.table" -> "table")).kudu //good 
   val kuduContext = new KuduContext("master:7051") //error
})

然后报错:

org.apache.spark.SparkException: 只有一个 SparkContext 可能正在运行 在这个 JVM 中(参见 SPARK-2243)。要忽略此错误,请设置 spark.driver.allowMultipleContexts = true。当前运行的 SparkContext 创建于: org.apache.spark.sql.SparkSession$Builder.getOrCreate(SparkSession.scala:860)

【问题讨论】:

  • 您似乎已经有一个活动的 SparkContext(因为您从 rdd.sparkContext.getConf 获得了配置。为什么要创建一个新的?
  • 我在 spark2-shell 中运行代码,默认包含 sparksession。
  • 如果您使用 spark-shell,则不需要 maven 依赖项。启动 shell 时包含 kudu jar。
  • 我可能会误导你。我现在更新了我的问题。
  • 你应该停止为每个 RDD 制作/获取新的 SparkSession 和 KuduContext

标签: scala apache-spark apache-kudu


【解决方案1】:

将您的 Kudu 版本更新到最新版本(当前为 1.5.0)。 KuduContext 在以后的版本中将SparkContext 作为输入参数,应该可以防止这个问题。

另外,在foreachRDD 之外进行初始 Spark 初始化。在您提供的代码中,将 sparkkuduContext 移出 foreach。此外,您不需要创建单独的sparkConf,您可以只使用较新的SparkSession

val spark = SparkSession.builder.appName("DirectKafka").master("local[*]").getOrCreate()
import spark.implicits._

val kuduContext = new KuduContext("master:7051", spark.sparkContext)
val bb = spark.read.options(Map("kudu.master" -> "master:7051", "kudu.table" -> "table")).kudu

val messages = KafkaUtils.createDirectStream("")
messages.foreachRDD(rdd => {   
  // do something with the bb table and messages       
})

【讨论】:

  • @cricket_007。用 kudu-spark2_2.11_1.1.0 ,好像只有一个参数 KuduContext(org.apache.kudu.spark.kudu)
  • 由于 spark 流文档,在 foreachRDD 内部进行了 Spark 初始化。出 foreachRD 有 val ssc = new StreamingContext(sparkConf, Seconds(2).
  • @Autumn:永远不需要在 foreach 中进行这种初始化。你在哪里看到的?
  • @Autumn:查看文档中链接的source code,他们实际上定义了一个在循环中使用的SparkSessionSingleton 对象。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-05-28
  • 2014-12-20
  • 1970-01-01
相关资源
最近更新 更多