【问题标题】:Scala Spark streaming kafkaScala Spark 流式卡夫卡
【发布时间】:2018-12-01 16:58:44
【问题描述】:

我在 kafka 中创建了一个示例主题,我正在尝试使用以下脚本在 spark 中使用内容:

import org.apache.spark._
 import org.apache.spark.streaming._
 import org.apache.spark.streaming.kafka._
 import org.apache.kafka.common.serialization.StringDeserializer
 import org.apache.spark.streaming.kafka010.LocationStrategies.PreferConsistent
import 
org.apache.spark.streaming.kafka010.ConsumerStrategies.Subscribe


 class Kafkaconsumer {
  val kafkaParams = Map[String, Object](
  "bootstrap.servers" -> "host1:port,host2:port2,host3:port3",
  "key.deserializer" -> classOf[StringDeserializer],
  "value.deserializer" -> classOf[StringDeserializer],
  "group.id" -> "use_a_separate_group_id_for_each_stream",
  "auto.offset.reset" -> "latest",
  "enable.auto.commit" -> (false: java.lang.Boolean)
  )
  val sparkConf = new SparkConf().setMaster("yarn")
 .setAppName("kafka example")
  val streamingContext = new StreamingContext(sparkConf, Seconds(10))
  val topics = Array("topicname")
  val topicsSet = topics.split(",").toSet
  val stream = KafkaUtils.createDirectStream[String, String](
  streamingContext,
  PreferConsistent,
  Subscribe[String, String](kafkaParams,topicsSet)
  )
  stream.print()
  stream.map(record => (record.key, record.value))
  streamingContext.start()
  streamingContext.awaitTermination()

我还包含了执行代码所需的库。

我有以下错误,请告诉我如何解决这个问题。

Error:
 Error:(23, 27) wrong number of type parameters for overloaded method value createDirectStream with alternatives:
  [K, V, KD <: kafka.serializer.Decoder[K], VD <: kafka.serializer.Decoder[V]](jssc: org.apache.spark.streaming.api.java.JavaStreamingContext, keyClass: Class[K], valueClass: Class[V], keyDecoderClass: Class[KD], valueDecoderClass: Class[VD], kafkaParams: java.util.Map[String,String], topics: java.util.Set[String])org.apache.spark.streaming.api.java.JavaPairInputDStream[K,V] <and>
  [K, V, KD <: kafka.serializer.Decoder[K], VD <: kafka.serializer.Decoder[V], R](jssc: org.apache.spark.streaming.api.java.JavaStreamingContext, keyClass: Class[K], valueClass: Class[V], keyDecoderClass: Class[KD], valueDecoderClass: Class[VD], recordClass: Class[R], kafkaParams: java.util.Map[String,String], fromOffsets: java.util.Map[kafka.common.TopicAndPartition,Long], messageHandler: org.apache.spark.api.java.function.Function[kafka.message.MessageAndMetadata[K,V],R])org.apache.spark.streaming.api.java.JavaInputDStream[R] <and>
  [K, V, KD <: kafka.serializer.Decoder[K], VD <: kafka.serializer.Decoder[V]](ssc: org.apache.spark.streaming.StreamingContext, kafkaParams: Map[String,String], topics: Set[String])(implicit evidence$19: scala.reflect.ClassTag[K], implicit evidence$20: scala.reflect.ClassTag[V], implicit evidence$21: scala.reflect.ClassTag[KD], implicit evidence$22: scala.reflect.ClassTag[VD])org.apache.spark.streaming.dstream.InputDStream[(K, V)] <and>
  [K, V, KD <: kafka.serializer.Decoder[K], VD <: kafka.serializer.Decoder[V], R](ssc: org.apache.spark.streaming.StreamingContext, kafkaParams: Map[String,String], fromOffsets: Map[kafka.common.TopicAndPartition,Long], messageHandler: kafka.message.MessageAndMetadata[K,V] => R)(implicit evidence$14: scala.reflect.ClassTag[K], implicit evidence$15: scala.reflect.ClassTag[V], implicit evidence$16: scala.reflect.ClassTag[KD], implicit evidence$17: scala.reflect.ClassTag[VD], implicit evidence$18: scala.reflect.ClassTag[R])org.apache.spark.streaming.dstream.InputDStream[R]val stream = KafkaUtils.createDirectStream[String, String](

【问题讨论】:

  • 能否请您发布整个错误跟踪?另外,这两条语句真的有效吗? val topics = Array("topicname") val topicsSet = topics.split(",").toSet - 因为您似乎试图在逗号符号上拆分数组。您可能会收到错误消息“split 不是 Array[String] 的成员”。
  • 感谢您的回复,因为您说我遇到了同样的错误。
  • 我出于其他原因添加了抱歉,现在我删除了该行并具有以下错误的相同代码:错误:(25、5)类型不匹配;发现:org.apache.spark.streaming.kafka010.LocationStrategy 需要:Map[String,String] PreferConsistent,
  • createDirectStream 语法总是有错误,当我添加 KafkaUtils.createDirectStream[String, String] 然后错误:重载方法值 createDirectStream 的类型参数数量错误。
  • 能否请您编辑帖子并分享您的各种尝试的详细信息?正如您在尝试中提到的两个不同问题,即“类型参数数量错误”和“类型不匹配”。我实际上是在尝试在我的系统中重现此问题,然后尝试修复。

标签: apache-spark kafka-consumer-api spark-streaming-kafka


【解决方案1】:

为你需要的key和value的解码类型添加类型参数,例如:

变化:

KafkaUtils.createDirectStream[String, String](
  streamingContext,
  PreferConsistent,
  Subscribe[String, String](kafkaParams,topicsSet)
)

KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder](
  streamingContext,
  PreferConsistent,
  Subscribe[String, String](kafkaParams,topicsSet)
)

【讨论】:

  • 非常感谢您的回复我被困在这一点上。
  • 你用的是什么版本的 spark、spark-streaming 和 spark-streaming-kafka,请分享一下
  • sparkVersion = "1.6.3",scalaVersion := "2.12.6",spark-streaming_2.10
  • 另外,如果您能给我提供有关如何从中构建胖罐子的见解,那对我真的很有帮助
  • 如果你使用spark-streaming_2.10 ,那么scala版本也应该是scalaVersion := "2.10.x
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-02-08
  • 2018-03-06
  • 2016-08-03
  • 2018-09-15
  • 2018-03-07
相关资源
最近更新 更多