diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java index b8be4bc081e8f..61a687754dcbc 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StoreChangelogReader.java @@ -41,6 +41,7 @@ import java.util.Set; public class StoreChangelogReader implements ChangelogReader { + private static final long RESTORE_LOG_INTERVAL_MS = 10_000L; private final Logger log; private final Consumer restoreConsumer; @@ -53,6 +54,8 @@ public class StoreChangelogReader implements ChangelogReader { private final Set completedRestorers = new HashSet<>(); private final Duration pollTime; + private long lastRestoreLogTime = 0L; + public StoreChangelogReader(final Consumer restoreConsumer, final Duration pollTime, final StateRestoreListener userStateRestoreListener, @@ -120,9 +123,42 @@ public Collection restore(final RestoringTasks active) { checkForCompletedRestoration(); + maybeLogRestorationProgress(); + return completedRestorers; } + private void maybeLogRestorationProgress() { + if (needsRestoring.isEmpty()) { + lastRestoreLogTime = 0L; + } else { + final long now = System.currentTimeMillis(); + if (now - lastRestoreLogTime > RESTORE_LOG_INTERVAL_MS) { + final Set topicPartitions = needsRestoring; + if (!topicPartitions.isEmpty()) { + final StringBuilder builder = new StringBuilder().append("Restoration in progress for ") + .append(topicPartitions.size()) + .append(" partitions."); + for (final TopicPartition partition : topicPartitions) { + final StateRestorer stateRestorer = stateRestorers.get(partition); + builder.append(" {") + .append(partition) + .append(": ") + .append("position=") + .append(stateRestorer.restoredOffset()) + .append(", end=") + .append(restoreToOffsets.get(partition)) + .append(", totalRestored=") + .append(stateRestorer.restoredNumRecords()) + .append("}"); + } + log.info(builder.toString()); + lastRestoreLogTime = now; + } + } + } + } + private void initialize(final RestoringTasks active) { if (!restoreConsumer.subscription().isEmpty()) { throw new StreamsException("Restore consumer should not be subscribed to any topics (" + restoreConsumer.subscription() + ")");