Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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<byte[], byte[]> restoreConsumer;
Expand All @@ -53,6 +54,8 @@ public class StoreChangelogReader implements ChangelogReader {
private final Set<TopicPartition> completedRestorers = new HashSet<>();
private final Duration pollTime;

private long lastRestoreLogTime = 0L;

public StoreChangelogReader(final Consumer<byte[], byte[]> restoreConsumer,
final Duration pollTime,
final StateRestoreListener userStateRestoreListener,
Expand Down Expand Up @@ -120,9 +123,42 @@ public Collection<TopicPartition> 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<TopicPartition> 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() + ")");
Expand Down