cadonna commented on code in PR #12279:
URL: https://github.com/apache/kafka/pull/12279#discussion_r899854007
##########
streams/src/main/java/org/apache/kafka/streams/processor/internals/DefaultStateUpdater.java:
##########
@@ -290,6 +295,34 @@ private void addTaskToRestoredTasks(final StreamTask task)
{
restoredActiveTasksLock.unlock();
}
}
+
+ private void maybeCommitRestoringTasks(final long now) {
+ final long elapsedMsSinceLastCommit = now - lastCommitMs;
+ if (elapsedMsSinceLastCommit > commitIntervalMs) {
+ if (log.isDebugEnabled()) {
+ log.debug("Committing all restoring tasks since {}ms has
elapsed (commit interval is {}ms)",
+ elapsedMsSinceLastCommit, commitIntervalMs);
+ }
+
+ for (final Task task : updatingTasks.values()) {
+ // do not enforce checkpointing during restoration if its
position has not advanced much
Review Comment:
I see! This is a bit hard to understand in my opinion. Could we have two
methods -- `commitTaskAndEnforceCheckpoint()` and
`commitTaskAndMaybeEnforcedCheckpoint()`? If we change the code to only write
the checkpoints, this code might change anyways.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]