【发布时间】:2020-02-07 23:06:20
【问题描述】:
我正在尝试执行以下代码,该代码从 Kafka Producer 提取消息并计算字数。
此代码来自 Github。 https://github.com/apache/spark/blob/master/examples/src/main/scala/org/apache/spark/examples/streaming/KafkaWordCount.scala
object KafkaWordCount {
def main(args: Array[String]) {
if (args.length < 4) {
System.err.println("Usage: KafkaWordCount <zkQuorum> <group> <topics> <numThreads>")
System.exit(1)
}
StreamingExamples.setStreamingLogLevels()
val Array(zkQuorum, group, topics, numThreads) = args
val sparkConf = new SparkConf().setAppName("KafkaWordCount")
val ssc = new StreamingContext(sparkConf, Seconds(2))
ssc.checkpoint("checkpoint")
val topicMap = topics.split(",").map((_, numThreads.toInt)).toMap
val lines = KafkaUtils.createStream(ssc, zkQuorum, group, topicMap).map(_._2)
val words = lines.flatMap(_.split(" "))
val wordCounts = words.map(x => (x, 1L))
.reduceByKeyAndWindow(_ + _, _ - _, Minutes(10), Seconds(2), 2)
wordCounts.print()
ssc.start()
ssc.awaitTermination()
}
}
我的问题是:我已经启动了 kafka zookeeper 和服务器(代理),并且还创建了一个名为“Mytopic1”的主题。 代码需要这 4 个参数:zkQuorum、group、topics、numThreads。 我知道我可以将“Mytopic1”、“1”作为主题和 numthreads 的参数传递。但是对于 zkQuorum 和 group 传递给我的应用程序的值是什么?
【问题讨论】:
标签: apache-spark apache-kafka spark-streaming