【问题标题】:Spark-shell Error object map is not a member of package org.apache.spark.streaming.rddSpark-shell 错误对象映射不是包 org.apache.spark.streaming.rdd 的成员
【发布时间】:2018-06-09 14:25:39
【问题描述】:

我正在尝试使用 spark 流从 Kafka 主题 KafkaStreamTestTopic1 读取 json 和解析两个值 valueStr1valueStr2。并将其转换为数据框以供进一步处理。

我在 spark-shell 中运行代码,因此可以使用 spark 上下文 sc

但是当我运行这个脚本时,它给了我以下错误:

错误:对象映射不是包 org.apache.spark.streaming.rdd 的成员 val dfa = rdd.map(record => {

下面是使用的脚本:

import org.apache.kafka.clients.consumer.ConsumerRecord
import org.apache.spark.{SparkConf, TaskContext}
import org.apache.spark._
import org.apache.spark.streaming._
import org.apache.spark.streaming.kafka010._
import org.apache.kafka.common.serialization.StringDeserializer
import play.api.libs.json._
import org.apache.spark.sql._

val ssc = new StreamingContext(sc, Seconds(5))

val sparkSession = SparkSession.builder().appName("myApp").getOrCreate()
val sqlContext = new SQLContext(sc)

// Create direct kafka stream with brokers and topics
val topicsSet = Array("KafkaStreamTestTopic1").toSet

// Set kafka Parameters
val kafkaParams = Map[String, String](
  "bootstrap.servers" -> "localhost:9092",
  "key.deserializer" -> "org.apache.kafka.common.serialization.StringDeserializer",
  "value.deserializer" -> "org.apache.kafka.common.serialization.StringDeserializer",
  "group.id" -> "my_group",
  "auto.offset.reset" -> "earliest",
  "enable.auto.commit" -> "false"
)

val stream = KafkaUtils.createDirectStream[String, String](
  ssc,
  LocationStrategies.PreferConsistent,
  ConsumerStrategies.Subscribe[String, String](topicsSet, kafkaParams)
)

val lines = stream.map(_.value)

lines.print()

case class MyObj(val one: JsValue)

lines.foreachRDD(rdd => {
  println("Debug Entered")

  import sparkSession.implicits._
  import sqlContext.implicits._


  val dfa = rdd.map(record => {

    implicit val myObjEncoder = org.apache.spark.sql.Encoders.kryo[MyObj]

    val json: JsValue = Json.parse(record)
    val value1 = (json \ "root" \ "child1" \ "child2" \ "valueStr1").getOrElse(null)
    val value2 = (json \ "root" \ "child1" \ "child2" \ "valueStr2").getOrElse(null)

    (new MyObj(value1), new MyObj(value2))

  }).toDF()

  dfa.show()
  println("Dfa Size is: " + dfa.count())


})

ssc.start()

【问题讨论】:

  • 您是否尝试过重命名您的rdd 以避免包冲突?

标签: scala apache-spark apache-spark-sql spark-streaming


【解决方案1】:

我想问题是 rdd 也是一个包(org.apache.spark.streaming.rdd),您使用以下行自动导入:

import org.apache.spark.streaming._

为避免此类冲突,请将您的变量重命名为其他名称,例如 myRdd

lines.foreachRDD(myRdd => { /* ... */ })

【讨论】:

  • 我现在收到一个新错误`error: value toDF is not a member of org.apache.spark.rdd.RDD[(MyObj, MyObj)]` at }).toDF()。请帮忙。
  • 你有没有在主函数的对象之外声明你的case类?
  • 我没有任何主要功能。只是我上面的脚本中显示的 shell 命令。在上面给定的脚本中,我在lines.foreachRDD(rdd => { 语句之前将我的案例类定义为case class MyObj(val one: JsValue)
  • 你能举个例子吗?
【解决方案2】:

将 spark-streaming 的依赖添加到构建管理器中

     "org.apache.spark" %% "spark-mllib" % SparkVersion,
    "org.apache.spark" %% "spark-streaming-kafka-0-10" % 
     "2.0.1"

您可以在构建过程中使用 maven 或 SBT 添加。

【讨论】:

  • 我没有使用任何 spark-mllib 函数
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-10-27
  • 2020-12-25
  • 2016-02-29
  • 1970-01-01
  • 2017-03-09
  • 2016-05-01
相关资源
最近更新 更多