【问题标题】:Spark Streaming Not reading all Kafka RecordsSpark Streaming 未读取所有 Kafka 记录
【发布时间】:2018-01-17 04:42:36
【问题描述】:

我们从 kafka 向 SparkStreaming 发送 15 条记录,但 spark 只接收 11 条记录。我正在使用 spark 2.1.0 和 kafka_2.12-0.10.2.0。

代码

import java.util.HashMap;
import java.util.Map;
import java.util.Properties;

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.spark.api.java.JavaSparkContext;
import org.apache.spark.api.java.function.Function;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.streaming.Duration;
import org.apache.spark.streaming.api.java.JavaDStream;
import org.apache.spark.streaming.api.java.JavaPairReceiverInputDStream;
import org.apache.spark.streaming.api.java.JavaStreamingContext;
import org.apache.spark.streaming.kafka.KafkaUtils;

import scala.Tuple2;

public class KafkaToSparkData {
 public static void main(String[] args) throws InterruptedException {
    int timeDuration = 100;
    int consumerNumberOfThreads = 1;
    String consumerTopic = "InputDataTopic";
    String zookeeperUrl = "localhost:2181";
    String consumerTopicGroup =  "testgroup";
    String producerKafkaUrl = "localhost:9092";
    String producerTopic =  "OutputDataTopic";
    String sparkMasterUrl = "local[2]";

    Map<String, Integer> topicMap = new HashMap<String, Integer>();
    topicMap.put(consumerTopic, consumerNumberOfThreads);

    SparkSession sparkSession = SparkSession.builder().master(sparkMasterUrl).appName("Kafka-Spark").getOrCreate();

    JavaSparkContext javaSparkContext = new JavaSparkContext(sparkSession.sparkContext());

    JavaStreamingContext javaStreamingContext = new JavaStreamingContext(javaSparkContext, new Duration(timeDuration));

    JavaPairReceiverInputDStream<String, String> messages = KafkaUtils.createStream(javaStreamingContext, zookeeperUrl, consumerTopicGroup, topicMap);

    JavaDStream<String> NewRecord = messages.map(new Function<Tuple2<String, String>, String>() {
        private static final long serialVersionUID = 1L;

        public String call(Tuple2<String, String> line) throws Exception {

            String responseToKafka = "";
            System.out.println(" Data IS " + line);

            String ValueData = line._2;
            responseToKafka = ValueData + "|" + "0";

            Properties configProperties = new Properties();
            configProperties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, producerKafkaUrl);
            configProperties.put("key.serializer", org.apache.kafka.common.serialization.StringSerializer.class);
            configProperties.put("value.serializer", org.apache.kafka.common.serialization.StringSerializer.class);

            KafkaProducer<String, String> producer = new KafkaProducer<String, String>(configProperties);

            ProducerRecord<String, String> topicMessage = new ProducerRecord<String, String>(producerTopic,responseToKafka);
            producer.send(topicMessage);
            producer.close();

            return responseToKafka;
        }
    });

    System.out.println(" Printing Record" );
    NewRecord.print();

    javaStreamingContext.start();
    javaStreamingContext.awaitTermination();
    javaStreamingContext.close();

    }
}
卡夫卡制作人

bin/kafka-console-producer.sh --broker-list localhost:9092 --topic InputDataTopic # 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18

卡夫卡消费者

bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic OutputDataTopic --from-beginning # 1|0 2|0 3|0 4|0 5|0 6|0 7|0 8|0 9|0 10|0 11|0

有人可以帮我解决这个问题吗?

【问题讨论】:

  • 您能添加一个完整的可重现示例吗? IE。还要添加生产者。
  • @maasg,我已经添加了完整的代码。我从 kafka 生产者发送了 18 条记录。但 kafka-consumer 只收到 11 条已处理的记录
  • 嗨,有人可以帮我一下吗。

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


【解决方案1】:

我们在这里看到的是惰性操作在 Spark 中的工作方式的影响。 在这里,我们使用map 操作来产生副作用,即向Kafka 发送一些数据。

然后使用print 实现流。默认情况下,print 将显示流的前 10 个元素,但采用 n+1 元素以显示“...”以指示何时有更多元素。

take(11) 强制实现前 11 个元素,因此它们从原始流中获取并使用 map 函数进行处理。这会导致部分发布到 Kafka。

如何解决这个问题?好吧,提示已经在上面了:不要在map 函数中使用副作用。 在这种情况下,消费流并将其发送到 Kafka 的正确输出操作应该是foreachRDD。

此外,为了避免为每个元素创建一个 Kafka 生产者实例,我们使用 foreachPartition 处理内部的 RDD。

这个过程的代码框架如下所示:

messages.foreachRDD{rdd => 
  rdd.foreachPartition{partitionIter => 
       producer = // create producer
       partitionIter.foreach{elem =>
           record = createRecord(elem)
           producer.send(record)
       }
       producer.flush()  
       producer.close()
   }
}    

【讨论】:

  • 来到这里发布我的解决方案(即在将 NewRecord.print(); 更改为 NewRecord.print(20); 后发现问题在于 print() 方法。因此将代码更改为 NewRecord.count( ).print();) 但现在从您那里找到了很好的解释和最佳方法/解决方案。感谢您的答复。现在我将遵循您的解决方案(即 foreachRDD)。
猜你喜欢
  • 2019-06-24
  • 2015-04-27
  • 2019-04-12
  • 2017-04-11
  • 2017-01-11
  • 2020-06-16
  • 1970-01-01
  • 2019-08-08
  • 2022-11-24
相关资源
最近更新 更多