【问题标题】:Kafka stream to spark does not reduce counts卡夫卡流火花不会减少计数
【发布时间】: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


    【解决方案1】:

    在阅读了几次代码后,我设法确定了为什么我总是将每个单词的 count 设为 1,而不是 总计

    在下面一行:

    // Create the context with 2 seconds batch size
    JavaStreamingContext jssc = new JavaStreamingContext(sparkConf, new Duration(2000));
    

    我将连续读取流的间隔设置为2秒。我意识到我没有在生产者端的这个间隔 (2 秒) 内产生足够的相同字符串的内容来获得聚合结果。

    但是,当我将此间隔增加到 10000 毫秒 (10 秒) 时,我可以从 kafka-producer。这些行由作业适当处理,并且在该特定时间间隔内很好地汇总了类似的字符串计数。

    (你好,4)

    (世界,6)

    非常感谢 沙比尔

    【讨论】:

      猜你喜欢
      • 2016-08-03
      • 2018-09-15
      • 2018-08-13
      • 2018-02-24
      • 2023-03-19
      • 2018-08-15
      • 1970-01-01
      • 2019-04-11
      • 1970-01-01
      相关资源
      最近更新 更多