【发布时间】:2016-12-19 12:34:55
【问题描述】:
我正在编写一个简单的 kafka - 在 eclipse 中使用 spark 流式传输代码来使用来自 kafka 代理的消息。下面是代码,当我尝试从 eclipse 运行代码时收到错误。
我还确保了依赖 jar 就位,请帮助摆脱这个错误
对象 spark_kafka_streaming {
def main(args: Array[String]) {
val conf = new SparkConf()
.setAppName("The swankiest Spark app ever")
.setMaster("local[*]")
val ssc = new StreamingContext(conf, Seconds(60))
ssc.checkpoint("C:\\keerthi\\software\\eclipse-jee-mars-2-win32- x86_64\\eclipse")
println("Parameters:" + "zkorum:" + "group:" + "topicMap:"+"number of threads:")
val zk = "xxxxxxxx:2181"
val group = "test-consumer-group"
val topics = "my-replicated-topic"
val numThreads = 2
val topicMap = topics.split(",").map((_,numThreads.toInt)).toMap
val lines = KafkaUtils.createStream(ssc,zk,group,topicMap).map(_._2)
val words = lines.flatMap(_.split(" "))
val wordCounts = words.map(x => (x,1L)).count()
println("wordCounts:"+wordCounts)
//wordCounts.print
}
}
例外:
线程“主”java.lang.NoClassDefFoundError 中的异常:org/apache/spark/streaming/kafka/KafkaUtils$ 在 org.firststream.spark_kakfa.spark_kafka_streaming$.main(spark_kafka_streaming.scala:30) 在 org.firststream.spark_kakfa.spark_kafka_streaming.main(spark_kafka_streaming.scala) 引起:java.lang.ClassNotFoundException:org.apache.spark.streaming.kafka.KafkaUtils$ 在 java.net.URLClassLoader.findClass(未知来源) 在 java.lang.ClassLoader.loadClass(未知来源) 在 sun.misc.Launcher$AppClassLoader.loadClass(未知来源) 在 java.lang.ClassLoader.loadClass(未知来源) ... 2 更多
依赖关系:
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka_2.10</artifactId>
<version>0.8.1.1</version>
<scope>compile</scope>
<exclusions>
<exclusion>
<artifactId>jmxri</artifactId>
<groupId>com.sun.jmx</groupId>
</exclusion>
<exclusion>
<artifactId>jms</artifactId>
<groupId>javax.jms</groupId>
</exclusion>
<exclusion>
<artifactId>jmxtools</artifactId>
<groupId>com.sun.jdmk</groupId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>0.8.2.0</version>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-streaming-kafka_2.10</artifactId>
<version>1.2.0</version>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-streaming_2.10</artifactId>
<version>1.2.0</version>
</dependency>
【问题讨论】:
-
下面是代码包org.firststream.spark_kakfa中用到的import kafka.serializer.StringDecoder import org.apache.spark.{ SparkContext, SparkConf } import org.apache.spark.streaming.{秒,StreamingContext } import org.apache.spark.streaming.kafka.KafkaUtils import org.apache.spark.streaming.kafka.KafkaUtils._ import org.apache.spark.streaming.kafka._
-
您的工作进展如何?你在制作超级 JAR 吗?
-
我从 eclipse 运行它,右键单击文件,通过 eclipse 作为 scala 应用程序运行
标签: apache-spark apache-kafka streaming spark-streaming-kafka