【问题标题】:How to make spark streaming commit in each batch when limiting Kafka batch size?限制Kafka批量大小时如何在每批中进行火花流提交?
【发布时间】:2020-08-19 01:16:20
【问题描述】:

为了在使用 Spark 流时限制批处理大小,我引用了这个answer

Kafka 中有大约 5000 万条记录存储(即将被消耗)。 主题有 3 个分区。

zhihu_comment   0          10906153        28668062        17761909        -               -               -
zhihu_comment   1          10972464        30271728        19299264        -               -               -
zhihu_comment   2          10906395        28662007        17755612        -               -               -

我的消费者应用:

public final class SparkConsumer {
  private static final Pattern SPACE = Pattern.compile(" ");

  public static void main(String[] args) throws Exception {
    String brokers = "device1:9092,device2:9092,device3:9092";
    String groupId = "spark";
    String topics = "zhihu_comment";

    // Create context with a certain seconds batch interval
    SparkConf sparkConf = new SparkConf().setAppName("TestKafkaStreaming");
    sparkConf.set("spark.streaming.backpressure.enabled", "true");
    sparkConf.set("spark.streaming.backpressure.initialRate", "10000");
    sparkConf.set("spark.streaming.kafka.maxRatePerPartition", "10000");
    JavaStreamingContext jssc = new JavaStreamingContext(sparkConf, Durations.seconds(10));

    Set<String> topicsSet = new HashSet<>(Arrays.asList(topics.split(",")));
    Map<String, Object> kafkaParams = new HashMap<>();
    kafkaParams.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, brokers);
    kafkaParams.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
    kafkaParams.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    kafkaParams.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);

    kafkaParams.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
    kafkaParams.put("enable.auto.commit", true);
    kafkaParams.put("max.poll.records", "500");


    // Create direct kafka stream with brokers and topics
    JavaInputDStream<ConsumerRecord<String, String>> messages = KafkaUtils.createDirectStream(
            jssc,
        LocationStrategies.PreferConsistent(),
        ConsumerStrategies.Subscribe(topicsSet, kafkaParams));

    // Get the lines, split them into words, count the words and print
    JavaDStream<String> lines = messages.map(ConsumerRecord::value);
    lines.count().print();

    jssc.start();
    jssc.awaitTermination();
  }
}

我已经限制了 spark 流的消耗大小,在我的例子中,我将 maxRatePerPartition 设置为 10000,这意味着在我的例子中它每批消耗 300000 条记录。

问题是虽然火花流能够处理具有特定限制的记录,the current offset showing by kafka is not the offset that spark streaming is handling. As the kafka's current offset suddenly goes down to latest offset:

zhihu_comment   0          28700537        28700676        139             consumer-1-ddcb0abd-e206-470d-925a-63ca4dc1d62a /192.168.0.102  consumer-1
zhihu_comment   1          30305102        30305224        122             consumer-1-ddcb0abd-e206-470d-925a-63ca4dc1d62a /192.168.0.102  consumer-1
zhihu_comment   2          28695033        28695146        113             consumer-1-ddcb0abd-e206-470d-925a-63ca4dc1d62a /192.168.0.102  consumer-1

看来Spark Streaming并没有在每一个batch中提交offset,它在开始消费的时候提交了最开始的offset!

有什么方法可以让每个批次的火花流提交?

Spark 流式日志,证明它每批消耗的记录数:

20/05/04 22:28:13 INFO scheduler.DAGScheduler: Job 15 finished: print at SparkConsumer.java:65, took 0.012606 s
-------------------------------------------
Time: 1588602490000 ms
-------------------------------------------
300000

20/05/04 22:28:13 INFO scheduler.JobScheduler: Finished job streaming job 1588602490000 ms.0 from job set of time 1588602490000 ms

【问题讨论】:

    标签: java apache-spark apache-kafka


    【解决方案1】:

    你需要禁用

    kafkaParams.put("enable.auto.commit", false);
    

    宁可使用

    messages.foreachRDD(rdd -> {
      OffsetRange[] offsetRanges = ((HasOffsetRanges) rdd.rdd()).offsetRanges();
    
      // do here some transformations and action on the rdd, typically like:
      // rdd.foreachPartition(it -> {
      //   it.foreach(row -> ...)
      // })
    
      // some time later, after outputs have completed
      ((CanCommitOffsets) messages.inputDStream()).commitAsync(offsetRanges);
    });
    

    Spark + Kafka Integration Guide 中所述。

    您也可以使用commitSync 进行同步提交。

    【讨论】:

    • 嗨@DennisLi,刚刚看到你的另一个post,看起来我上面的回答不够清楚。我在回答中给出的代码 sn-p 也可以调整为进行常规 rdd 处理。只需将所有操作和转换添加到编写评论的部分。给定的代码 sn-p 并不意味着是独立的。
    • 嗨,mike,谢谢你的回答,我试过后很有帮助。我不太清楚你的意思。你的意思是像“messages.flatMap(record -&gt; Arrays.asList(record.value()).iterator());”这样的sn-p可以添加到messages.foreachRDD中?
    • flatMap 应该在 foreachRDD 方法之外。但一般情况下,您不需要调用两次messages.foreachRDD
    • 查看您其他帖子的代码,您可以将我上面答案中给出的偏移部分放入recordStream.foreachRDD
    • 嗨,迈克,你能在你的答案中更新它吗?我试过了,但它引发了 "cannot be cast to org.apache.spark.streaming.kafka010.HasOffsetRanges" 。既然messages是JavaInputDStream,recordStream是JavaDStream,有区别吗?
    猜你喜欢
    • 2017-02-20
    • 2023-03-16
    • 2019-03-28
    • 2016-06-25
    • 2016-10-26
    • 2016-10-02
    • 1970-01-01
    • 2018-07-16
    • 1970-01-01
    相关资源
    最近更新 更多