【发布时间】:2018-08-15 21:52:02
【问题描述】:
我是新的 Spark 和 Kafka。我已经在 Windows 系统中完成了 Spark 和 Kafka 的设置。两者都工作得很好。我已经按照下面提到的教程并在 Spark-shell 中执行 Scala 代码并得到下面提到的错误。
谁能帮助我了解如何在 Scala 中使用 Spark 收听 Kafka 流。
火花:2.2 卡夫卡:2.12
http://www.godatafy.com/poc/streaming-with-spark-kafka/
spark-shell -jars ..\jars\spark-streaming-kafka_2.11-1.6.3.jar
scala> import org.apache.spark.SparkConf
import org.apache.spark.SparkConf
scala> import org.apache.spark.streaming.StreamingContext
import org.apache.spark.streaming.StreamingContext
scala> import org.apache.spark.streaming.Seconds
import org.apache.spark.streaming.Seconds
scala> import org.apache.spark.streaming.kafka.KafkaUtils
import org.apache.spark.streaming.kafka.KafkaUtils
scala> sc.stop
scala> val sparkConf = new SparkConf().setAppName("KafkaWordCount").setMaster("local[2]")
sparkConf: org.apache.spark.SparkConf = org.apache.spark.SparkConf@427c2c96
scala> val ssc = new StreamingContext(sparkConf, Seconds(2))
ssc: org.apache.spark.streaming.StreamingContext = org.apache.spark.streaming.StreamingContext@1fd73dcb
scala> val lines = KafkaUtils.createStream(ssc, "localhost:2181", "spark-streaming-consumer-group", Map("test" -> 5))
error: missing or invalid dependency detected while loading class file 'KafkaUtils.class'.
Could not access term kafka in package <root>,
because it (or its dependencies) are missing. Check your build definition for
missing or conflicting dependencies. (Re-run with `-Ylog-classpath` to see the problematic classpath.)
A full rebuild may help if 'KafkaUtils.class' was compiled against an incompatible version of <root>.
【问题讨论】:
-
你确定 kafka 2.12 吗?
-
是的。我正在使用 Kafka:2.12。
标签: scala apache-spark apache-kafka spark-streaming