【问题标题】:Kafka popular hashtags counting卡夫卡流行标签计数
【发布时间】:2018-02-17 01:59:12
【问题描述】:

我正在使用 Kafka 和 Spark 来统计 Twitter 上最受欢迎的主题标签。所以,这是我要运行的 scala 对象:

package spark.example

import java.util.HashMap

import org.apache.kafka.clients.producer.{ KafkaProducer, ProducerConfig, ProducerRecord }
import org.apache.spark.SparkConf
import org.apache.spark.streaming._
import org.apache.spark.streaming.kafka._
import org.apache.spark.streaming.{ Seconds, StreamingContext }
import org.apache.spark.SparkContext._
import org.apache.spark.streaming.twitter._
import org.apache.spark.SparkConf
import org.apache.spark.streaming._
import org.apache.spark.{ SparkContext, SparkConf }
import org.apache.spark.storage.StorageLevel

/**
  * A Spark Streaming - Kafka integration to receive twitter
  * data from kafka topic and find the popular hashtags
  *
  * Arguments: <zkQuorum> <consumer-group> <topics> <numThreads>
  * <zkQuorum>       - The zookeeper hostname
  * <consumer-group> - The Kafka consumer group
  * <topics>         - The kafka topic to subscribe to
  * <numThreads>     - Number of kafka receivers to run in parallel
  *
  * More discussion at stdatalabs.blogspot.com
  *
  * @author Sachin Thirumala
  */

object KafkaSparkPopularHashTags {

  val conf = new SparkConf().setMaster("local[6]").setAppName("Spark Streaming - Kafka Producer - PopularHashTags").set("spark.executor.memory", "1g")

  conf.set("spark.streaming.receiver.writeAheadLog.enable", "true")

  val sc = new SparkContext(conf)

  def main(args: Array[String]) {

//sc.setLogLevel("WARN")

// Create an array of arguments: zookeeper hostname/ip,consumer group, topicname, num of threads
val Array(zkQuorum, group, topics, numThreads) = args

// Set the Spark StreamingContext to create a DStream for every 2 seconds
val ssc = new StreamingContext(sc, Seconds(2))
ssc.checkpoint("checkpoint")

// Map each topic to a thread
val topicMap = topics.split(",").map((_, numThreads.toInt)).toMap
// Map value from the kafka message (k, v) pair
val lines = KafkaUtils.createStream(ssc, zkQuorum, group, topicMap).map(_._2)
// Filter hashtags
val hashTags = lines.flatMap(_.split(" ")).filter(_.startsWith("#"))

// Get the top hashtags over the previous 60/10 sec window
val topCounts60 = hashTags.map((_, 1)).reduceByKeyAndWindow(_ + _, Seconds(60))
  .map { case (topic, count) => (count, topic) }


val topCounts10 = hashTags.map((_, 1)).reduceByKeyAndWindow(_ + _, Seconds(10))
  .map { case (topic, count) => (count, topic) }


lines.print()

// Print popular hashtags
topCounts60.foreachRDD(rdd => {
  val topList = rdd.take(10)
  println("\nPopular topics in last 60 seconds (%s total):".format(rdd.count()))
  topList.foreach { case (count, tag) => println("%s (%s tweets)".format(tag, count)) }
})

topCounts10.foreachRDD(rdd => {
  val topList = rdd.take(10)
  println("\nPopular topics in last 10 seconds (%s total):".format(rdd.count()))
  topList.foreach { case (count, tag) => println("%s (%s tweets)".format(tag, count)) }
})

lines.count().map(cnt => "Received " + cnt + " kafka messages.").print()

ssc.start()
ssc.awaitTermination()
  }
}

但每次我尝试使用以下参数运行代码时: 本地主机:2181 spark-streaming-consumer-group 推文 2 (已创建主题“推文”且 TwitterProducer 正在运行)我收到以下错误:

Exception in thread "main" java.lang.InstantiationException: org.apache.spark.util.SystemClock
at java.lang.Class.newInstance(Class.java:427)
at org.apache.spark.streaming.scheduler.JobGenerator.liftedTree1$1(JobGenerator.scala:52)
at org.apache.spark.streaming.scheduler.JobGenerator.<init>(JobGenerator.scala:51)
at org.apache.spark.streaming.scheduler.JobScheduler.<init>(JobScheduler.scala:54)
at org.apache.spark.streaming.StreamingContext.<init>(StreamingContext.scala:183)
at org.apache.spark.streaming.StreamingContext.<init>(StreamingContext.scala:75)
at spark.example.KafkaSparkPopularHashTags$.main(KafkaSparkPopularHashTags.scala:48)
at spark.example.KafkaSparkPopularHashTags.main(KafkaSparkPopularHashTags.scala)
Caused by: java.lang.NoSuchMethodException: org.apache.spark.util.SystemClock.<init>()
at java.lang.Class.getConstructor0(Class.java:3082)
at java.lang.Class.newInstance(Class.java:412)

由于错误表明第 48 行有问题,即:

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

系统似乎无法实例化该对象。 有什么建议可以解决这个问题吗? 其他信息:我正在使用 Scala 2.12,并且我已经尝试降级到 Scala 2.11 甚至 2.10。 我正在尝试重现这个实验:http://stdatalabs.blogspot.in/2016/09/spark-streaming-part-3-real-time.html

【问题讨论】:

  • 建议检查依赖版本。
  • 同意这是一个依赖问题,而不是您向我们展示的代码中的问题。特定版本的 Spark 需要特定的主要版本 Scala(2.10、2.11 或(在某些尚未正式发布的 Spark 版本中)2.12)。为不同的 Scala 主要版本构建的库不能混合使用。 Spark 的当前稳定版本使用 Scala 2.11,因此您可能希望坚持使用 2.11 并确保您只使用 2.11 库。
  • 经过多次尝试,我发现我必须使用 Scala 版本 2.11.1 感谢您的宝贵帮助

标签: scala twitter apache-kafka


【解决方案1】:

[更新]经过多次尝试,我发现我必须使用 Scala 版本 2.11.1 感谢您的宝贵帮助

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-10-02
    • 2016-05-11
    • 2017-02-08
    • 2018-03-06
    • 2016-08-03
    相关资源
    最近更新 更多