【问题标题】:Duplicates while Processing Streaming Data using Kafka-Spark Streaming API使用 Kafka-Spark Streaming API 处理流数据时重复
【发布时间】:2019-06-28 15:02:56
【问题描述】:

以下代码在数据处理后工作并提交偏移量。 但问题是,它在以下情况下处理重复:

消费者作业正在运行,hive表有0条记录,当前偏移量为(FORMAT-fromOffest, untilOffset, Difference): 512 512 0

然后我生成了 1000 条记录,当它读取 34 条记录但未提交时,我将其杀死 512 546 34

我看到此时,34 个记录已经加载到 Hive 表中

接下来,我重新启动了应用程序。

我看到它再次读取了 34 条记录(而不是读取 1000-34=76 条记录),尽管它已经处理了它们并加载到 Hive 512 1512 1000 然后几秒钟后它会更新。 1512 1512 0 Hive 现在有 (34+1000=1034)

这会导致表中出现重复记录(额外 34 条)。 如代码中所述,我仅在处理/加载到 Hive 表后才提交偏移量。

public void method1(SparkConf conf,String app) 
    spark = SparkSession.builder().appName(conf.get("")).enableHiveSupport().getOrCreate();
    final JavaStreamingContext javaStreamContext = new JavaStreamingContext(context,
            new Duration(<spark duration>));
    JavaInputDStream<ConsumerRecord<String, String>> messages = KafkaUtils.createDirectStream(javaStreamContext,
            LocationStrategies.PreferConsistent(),
            ConsumerStrategies.<String, String> Subscribe(<topicnames>, <kafka Params>));

            JavaDStream<String> records = messages.map(new Function<ConsumerRecord<String, String>, String>() {
                @Override
                public String call(ConsumerRecord<String, String> tuple2) throws Exception {
                    return tuple2.value();
                }
            });

            records.foreachRDD(new VoidFunction<JavaRDD<String>>() {
                @Override
                public void call(JavaRDD<String> rdd) throws Exception {
                    if(!rdd.isEmpty()) {
                        methodToSaveDataInHive(rdd, <StructTypeSchema>,<OtherParams>);
                    }
                }
             });

             messages.foreachRDD(new VoidFunction<JavaRDD<ConsumerRecord<String, String>>>() {
              @Override
              public void call(JavaRDD<ConsumerRecord<String, String>> rdd) {
                    OffsetRange[] offsetRanges = ((HasOffsetRanges) rdd.rdd()).offsetRanges();
                    ((CanCommitOffsets) messages.inputDStream()).commitAsync(offsetRanges);                     
                    for (OffsetRange offset : offsetRanges) {
                        System.out.println(offset.fromOffset() + " " + offset.untilOffset()+ "  "+offset.count());
                    }
                     }
              });             
    javaStreamContext.start();
    javaStreamContext.awaitTermination();
}

【问题讨论】:

    标签: java apache-kafka spark-streaming offset kafka-consumer-api


    【解决方案1】:

    一般来说,在构建 Spark Streaming 作业时,您不应该担心重复,而应该在下游处理。不要误会我的意思,您希望构建应用程序以防止重复,但是当发生灾难性事情时,您会得到重复,这就是为什么最好稍后再管理它的原因。

    我看到的第一个问题是您在哪里保存偏移量。您应该在保存数据后立即保存它们,而不是之后的方法。当 records.foreachRDD 完成 methodToSaveData 时,它应该调用以保存偏移量。您可能需要重新构建映射记录的方式,以便获得偏移量详细信息,但这是最好的地方。

            records.foreachRDD(new VoidFunction<JavaRDD<String>>() {
                @Override
                public void call(JavaRDD<String> rdd) throws Exception {
                    if(!rdd.isEmpty()) {
                        methodToSaveDataInHive(rdd, <StructTypeSchema>,<OtherParams>);
                        **{commit offsets here}**
                    }
                }
             });
    

    也就是说,保存偏移量的位置并不重要。如果作业在将数据写入 hive 之后且在提交偏移范围之前被终止,您将重新处理记录。有一些方法可以构建应用程序,因此它具有优雅的关闭挂钩(Google it),它试图捕获一个 kill 命令并优雅地关闭它,但同样容易受到应用程序如何被杀死或崩溃的影响。如果运行 executor 的机器在保存到 hive 之后但在提交偏移量之前失去了电源,那么你就有了重复项。如果应用程序被终止 -9(在 Linux 中),它不关心正常关闭,您将有重复。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-10-10
      • 2016-03-12
      • 2017-02-06
      • 2016-03-04
      • 2017-07-13
      • 2019-10-06
      • 2021-05-22
      • 1970-01-01
      相关资源
      最近更新 更多