【发布时间】:2017-10-02 02:50:09
【问题描述】:
我正在尝试从 kafka 流式传输到 spark 的基本示例。我对 spark 很陌生,经验很少。
我的程序如下(复制自apache-spark中的例子):
if (args.length < 4) {
System.err.println("Usage: JavaKafkaWordCount <zkQuorum> <group> <topics> <numThreads>");
System.exit(1);
}
String zkQuorum = args[0];
String groupId = args[1];
String topicsToListen = args[2];
String numOfThread = args[3];
StreamingExamples.setStreamingLogLevels();
SparkConf sparkConf = new SparkConf().setAppName("JavaKafkaWordCount");
// Create the context with 2 seconds batch size
JavaStreamingContext jssc = new JavaStreamingContext(sparkConf, new Duration(2000));
int numThreads = Integer.parseInt(numOfThread);
Map<String, Integer> topicMap = new HashMap<>();
String[] topics = topicsToListen.split(",");
for (String topic : topics) {
topicMap.put(topic, numThreads);
}
JavaPairReceiverInputDStream<String, String> messages =
KafkaUtils.createStream(jssc, zkQuorum, groupId, topicMap);
JavaDStream<String> lines = messages.map(Tuple2::_2);
JavaDStream<String> words = lines.flatMap(x -> Arrays.asList(SPACE.split(x)).iterator());
JavaPairDStream<String, Integer> wordCounts = words.mapToPair(s -> new Tuple2<>(s, 1))
.reduceByKey((i1, i2) -> i1 + i2);
wordCounts.print();
jssc.start();
jssc.awaitTermination();
然后我启动我的 kafka-broker 并通过生成以下命令运行构建的 jar:
$SPARK_HOME/bin/spark-submit --class "JavaKafkaWordCount" --master local[2] PATH_TO_JAR/kafka-spark-streaming-1.0-SNAPSHOT-jar-with-dependencies.jar localhost:2181 test-consumer-小组测试1
当我从 kafka-producer 生成一些单词时,我预计多次发布的单词计数会增加,但我看到的只是单词和 count strong> 为每个新发布打印 1 个:
(你好,1)
当我多次发布同一个词时,我预计计数会增加,
(你好,2)
但这并没有发生。我在这里到底理解错了什么,这与我传递给他的工作的论点有关还是工作的本意是什么?
有人可以提供一些见解吗?
谢谢 沙比尔
【问题讨论】:
标签: apache-spark streaming apache-kafka