【问题标题】:Input arguments zkQuorum and group to be passed to Kafkautils.createStream method要传递给 Kafkautils.createStream 方法的输入参数 zkQuorum 和组
【发布时间】: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


    【解决方案1】:
    • zkQuorum 是一个逗号分隔的 Zookeper URI 列表,例如 your.host.name:2181
    • group 是您要用于您的消费者组的任意组名。

    【讨论】:

    • 如何创建消费组?我使用以下命令在 Windows 命令提示符下启动了消费者。那么我将使用什么消费者群体作为论据呢?命令:“kafka-console-consumer.bat --bootstrap-server localhost:9092 --topic Mytopic1 --from-beginning”
    • fyi 我能够通过命令提示符从主题“Mytopic1”发送消息,并且能够在我的命令提示符下接收消费者的消息。但正如我所说,我的要求是,我希望我的 Spark 应用程序通过传递所需的参数来读取来自 Kafka 主题的消息。卡在这里了。你能在这方面帮忙吗?
    猜你喜欢
    • 2013-01-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-03-16
    • 2017-07-01
    • 1970-01-01
    • 1970-01-01
    • 2021-03-16
    相关资源
    最近更新 更多