apache / apache/hudi

Fix NPE when there is mismatch in num of kafka partitions

Open
#16,080 0 comments 0 reactions 0 assignees View on GitHub
area:ingest from-jira priority:high type:bug
Dominant language
Java
Stars
6.2k
Forks
2.5k
Avg merge
2d 8h
Merged PRs (30d)
111

Description

lets say latest checkpoint in deltastreamer has 5 kafka partition offset. Users deleted and re-created and now there is only one kafka partition. So, when checking for new offsets, we run into NPE. 

We might need to fix the parsing logic to accommodate only new partitions. 

stacktrace:
{code:java}
23/07/06 16:07:52 ERROR DeltaStreamer : Failed to run job for table: ABC
java.lang.NullPointerException
at org.apache.hudi.utilities.sources.helpers.KafkaOffsetGen.lambda$fetchValidOffsets$1(KafkaOffsetGen.java:404)
at java.util.stream.MatchOps$1MatchSink.accept(MatchOps.java:90)
at java.util.HashMap$EntrySpliterator.tryAdvance(HashMap.java:1744)
at java.util.stream.ReferencePipeline.forEachWithCancel(ReferencePipeline.java:126)
at java.util.stream.AbstractPipeline.copyIntoWithCancel(AbstractPipeline.java:499)
at java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:486)
at java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:472)
at java.util.stream.MatchOps$MatchOp.evaluateSequential(MatchOps.java:230)
at java.util.stream.MatchOps$MatchOp.evaluateSequential(MatchOps.java:196)
at java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234)
at java.util.stream.ReferencePipeline.anyMatch(ReferencePipeline.java:516)
at org.apache.hudi.utilities.sources.helpers.KafkaOffsetGen.fetchValidOffsets(KafkaOffsetGen.java:404)
at org.apache.hudi.utilities.sources.helpers.KafkaOffsetGen.getNextOffsetRanges(KafkaOffsetGen.java:317)
at org.apache.hudi.utilities.sources.KafkaSource.fetchNewData(KafkaSource.java:67)
at org.apache.hudi.utilities.sources.Source.fetchNext(Source.java:105)
at org.apache.hudi.utilities.deltastreamer.SourceFormatAdapter.fetchNewDataInRowFormat(SourceFormatAdapter.java:288)
at org.apache.hudi.utilities.deltastreamer.DeltaSync.fetchFromSource(DeltaSync.java:477)
at org.apache.hudi.utilities.deltastreamer.DeltaSync.readFromSource(DeltaSync.java:451)
at org.apache.hudi.utilities.deltastreamer.DeltaSync.syncOnce(DeltaSync.java:358)
at org.apache.hudi.utilities.deltastreamer.HoodieDeltaStreamer$DeltaSyncService.ingestOnce(HoodieDeltaStreamer.java:876)
at org.apache.hudi.common.util.Option.ifPresent(Option.java:97)
at {code}
Code snippet of interest 
{code:java}
private Map fetchValidOffsets(KafkaConsumer consumer,
Option lastCheckpointStr, Set topicPartitions) {
Map earliestOffsets = consumer.beginningOffsets(topicPartitions);
Map checkpointOffsets = CheckpointUtils.strToOffsets(lastCheckpointStr.get());
boolean isCheckpointOutOfBounds = checkpointOffsets.entrySet().stream()
.anyMatch(offset -> offset.getValue() < earliestOffsets.get(offset.getKey())); {code}
last time where we do earliestOffsets.get(offset.getKey()) runs into NPE. 

 

## JIRA info

- Link: https://issues.apache.org/jira/browse/HUDI-6502
- Type: Bug

Contributor guide

No contributing guide indexed for this repository

Research direction

Start in org.apache.hudi.utilities.sources.helpers.KafkaOffsetGen.fetchValidOffsets at KafkaOffsetGen.java:404, then inspect CheckpointUtils.strToOffsets and the callers shown in the stack trace. Reproduce a checkpoint with five partitions against a topic with one, and verify the completed change handles missing partitions without an NPE.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
data-engineering, stream-processing
Issue type
Bug
Difficulty
3/5
Estimated time
1-2 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
48/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.