【问题标题】:Getting empty values while receiving from kafka Spark streaming从kafka Spark流中接收时获取空值
【发布时间】:2017-01-25 23:52:14
【问题描述】:

我是 Spark 流媒体 的新手,我正在实施小练习,例如从 kafka 发送 XML 数据,并且需要接收 >通过火花流传输数据。我尝试了所有可能的方式..但每次我得到空值。

Kafka端没有问题,唯一的问题是从Spark端接收流数据。

这是我如何实现的代码:

package com.package;

import org.apache.spark.SparkConf;
import org.apache.spark.api.java.JavaSparkContext;
import org.apache.spark.streaming.Duration;
import org.apache.spark.streaming.api.java.JavaStreamingContext;

public class SparkStringConsumer {

    public static void main(String[] args) {

        SparkConf conf = new SparkConf()
                .setAppName("kafka-sandbox")
                .setMaster("local[*]");
        JavaSparkContext sc = new JavaSparkContext(conf);
        JavaStreamingContext ssc = new JavaStreamingContext(sc, new Duration(2000));

        Map<String, String> kafkaParams = new HashMap<>();
        kafkaParams.put("metadata.broker.list", "localhost:9092");
        Set<String> topics = Collections.singleton("mytopic");

        JavaPairInputDStream<String, String> directKafkaStream = KafkaUtils.createDirectStream(ssc,
        String.class, String.class, StringDecoder.class, StringDecoder.class, kafkaParams, topics);
        directKafkaStream.foreachRDD(rdd -> {
        System.out.println("--- New RDD with " + rdd.partitions().size()
            + " partitions and " + rdd.count() + " records");
        rdd.foreach(record -> System.out.println(record._2));
        });


        ssc.start();
        ssc.awaitTermination();
    }
}

我正在使用以下版本:

**动物园管理员 3.4.6

Scala 2.11

火花 2.0

卡夫卡 0.8.2**

【问题讨论】:

    标签: hadoop apache-spark streaming apache-kafka spark-streaming


    【解决方案1】:

    你可以这样:

    directKafkaStream.foreachRDD(rdd ->{            
                rdd.foreachPartition(item ->{
                    while (item.hasNext()) {    
                        System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>"+item.next());
    }
    }
    });
    

    itme.next() 包含键值对。你可以通过使用来获取值 item.next()._2

    【讨论】:

      【解决方案2】:

      您的 spark 流应用程序看起来不错。我对其进行了测试,它正在打印 kafka 消息。您也可以尝试下面的“收到消息”打印语句来验证 kafka 消息。

          directKafkaStream.foreachRDD(rdd -> {
          System.out.println("Message Received "+rdd.values().take(5));
          System.out.println("--- New RDD with " + rdd.partitions().size()
              + " partitions and " + rdd.count() + " records");
          rdd.foreach(record -> System.out.println(record._2));
          });
      

      如果您使用的是 Zookeeper,那么也将其设置为 kafka 参数

      kafkaParams.put("zookeeper.connect","localhost:2181");
      

      在您的程序中没有看到以下导入语句,因此在此处添加。

      import org.apache.spark.streaming.kafka.KafkaUtils;
      import kafka.serializer.StringDecoder;
      

      还请验证您是否可以使用命令行 kafka-console-consumer 消费关于主题“mytopic”的消息。

      【讨论】:

      • 您好@abaghel,感谢您的快速回复。我按照你说的尝试过,仍然收到空消息...这是消息:收到的消息 [] --- 具有 1 个分区和 0 条记录的新 RDD
      • 你试过 kafka-console-consumer 吗?你能看到那里的消息吗?
      • 有趣。你能分享你的 pom.xml 吗?
      • 您能否分享您的电子邮件 ID,否则我将在此处粘贴我的 pom.xml 文件....或者这是我的邮件 ID contacteedupuganti@gmail.com 您可以在此邮件 ID 上留言。 ..
      猜你喜欢
      • 2017-04-07
      • 2021-02-21
      • 2017-05-09
      • 1970-01-01
      • 2019-02-17
      • 2020-04-11
      • 2018-02-05
      • 2016-03-04
      • 2017-01-26
      相关资源
      最近更新 更多