【发布时间】: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