NIFI-15293 Checkpoint committed records in ConsumeKinesis#10600
Merged
exceptionfactory merged 1 commit intoapache:mainfrom Dec 4, 2025
Conversation
awelless
commented
Dec 4, 2025
| void doCheckpoint() throws KinesisClientLibDependencyException, InvalidStateException, ThrottlingException, ShutdownException, IllegalArgumentException; | ||
| } | ||
|
|
||
| private record LastIgnoredCheckpoint(RecordProcessorCheckpointer checkpointer, String sequenceNumber, long subSequenceNumber) { |
Contributor
Author
There was a problem hiding this comment.
In the actual KCL implementation the same stateful RecordProcessorCheckpointer is passed to each processor's method, however the API doesn't call this out explicitly.
In case this changes in the future, it's better to keep a reference to the checkpointer which was passed with a batch of records.
906cfcc to
5cac719
Compare
exceptionfactory
approved these changes
Dec 4, 2025
Contributor
exceptionfactory
left a comment
There was a problem hiding this comment.
Thanks for addressing this issue @awelless, the updated approach looks good. +1 merging
mark-bathori
pushed a commit
to mark-bathori/nifi
that referenced
this pull request
Feb 5, 2026
Signed-off-by: David Handermann <exceptionfactory@apache.org>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
NIFI-15293
The pull request addresses potential data loss caused by incorrect checkpointing logic. See the jira ticket for more details.
Also tested manually, with abrupt interruptions and shutdowns of the processor.
Tracking
Please complete the following tracking steps prior to pull request creation.
Issue Tracking
Pull Request Tracking
NIFI-00000NIFI-00000Pull Request Formatting
mainbranchVerification
Please indicate the verification steps performed prior to pull request creation.
Build
./mvnw clean install -P contrib-checkLicensing
LICENSEandNOTICEfilesDocumentation