【发布时间】:2017-02-11 07:35:45
【问题描述】:
我正在使用 Spark Streaming 来使用来自 Kafka 主题的数据。
如果我使用DirectStream 方法,我没有定义consumer group 和number of consumers 的选项。
例如:
val messages = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder](ssc, kafkaParams, topicsSet)
我在哪里定义消费者组和该组的消费者数量?
如果我使用基于 Receiver 的方法,我可以选择定义 consumer group 和 number of threads[此组中的消费者数量]。
基于接收者的方法:
val topicMap = topics.split(",").map((_, numThreads.toInt)).toMap
val lines = KafkaUtils.createStream(ssc, zkQuorum, group, topicMap).map(_._2)
【问题讨论】:
标签: scala spark-streaming kafka-consumer-api