[
https://issues.apache.org/jira/browse/HUDI-1006?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Tianye Li updated HUDI-1006:
----------------------------
Summary: deltastreamer use kafkaSource with offset reset strategy: latest
can't consume data (was: deltastreamer use kafkaSource set
auto.offset.reset=latest can't consume data)
> deltastreamer use kafkaSource with offset reset strategy: latest can't
> consume data
> -----------------------------------------------------------------------------------
>
> Key: HUDI-1006
> URL: https://issues.apache.org/jira/browse/HUDI-1006
> Project: Apache Hudi
> Issue Type: Bug
> Components: DeltaStreamer
> Reporter: liujinhui
> Assignee: Tianye Li
> Priority: Major
> Fix For: 0.6.0
>
>
> org.apache.hudi.utilities.sources.JsonKafkaSource#fetchNewData
> if (totalNewMsgs <= 0) {
> return new InputBatch<>(Option.empty(), lastCheckpointStr.isPresent() ?
> lastCheckpointStr.get() : "");
> }
> I think it should not be empty here, it should be
> if (totalNewMsgs <= 0) {
> return new InputBatch<>(Option.empty(), lastCheckpointStr.isPresent() ?
> lastCheckpointStr.get() : CheckpointUtils.offsetsToStr(offsetRanges));
> }
--
This message was sent by Atlassian Jira
(v8.3.4#803005)