From cacf94f618fd659ff9b40f89a590a2285a55bb80 Mon Sep 17 00:00:00 2001 From: Bill Bejeck Date: Tue, 26 Sep 2017 18:33:50 -0400 Subject: [PATCH 1/4] KAFKA-5958: add StateRestoreListener to GlobalThread for restore progress --- .../org/apache/kafka/streams/KafkaStreams.java | 3 ++- .../internals/GlobalStateManagerImpl.java | 17 ++++++++++++++--- .../processor/internals/GlobalStreamThread.java | 14 +++++++++----- .../apache/kafka/streams/KafkaStreamsTest.java | 5 +++-- .../internals/GlobalStateManagerImplTest.java | 6 ++++-- .../internals/GlobalStreamThreadTest.java | 6 ++++-- .../kafka/test/ProcessorTopologyTestDriver.java | 3 ++- 7 files changed, 38 insertions(+), 16 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java b/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java index 5aec3c5a97171..2f5ce4bd37909 100644 --- a/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java +++ b/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java @@ -613,7 +613,8 @@ public void onRestoreEnd(final TopicPartition topicPartition, final String store stateDirectory, metrics, Time.SYSTEM, - globalThreadId); + globalThreadId, + delegatingStateRestoreListener); globalThreadState = globalStreamThread.state(); } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImpl.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImpl.java index d9205a0c4253b..d03425bf88147 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImpl.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImpl.java @@ -27,6 +27,7 @@ import org.apache.kafka.streams.errors.StreamsException; import org.apache.kafka.streams.processor.BatchingStateRestoreCallback; import org.apache.kafka.streams.processor.StateRestoreCallback; +import org.apache.kafka.streams.processor.StateRestoreListener; import org.apache.kafka.streams.processor.StateStore; import org.apache.kafka.streams.state.internals.OffsetCheckpoint; import org.slf4j.Logger; @@ -61,15 +62,18 @@ public class GlobalStateManagerImpl implements GlobalStateManager { private final OffsetCheckpoint checkpoint; private final Set globalStoreNames = new HashSet<>(); private final Map checkpointableOffsets = new HashMap<>(); + private final StateRestoreListener stateRestoreListener; public GlobalStateManagerImpl(final ProcessorTopology topology, final Consumer consumer, - final StateDirectory stateDirectory) { + final StateDirectory stateDirectory, + final StateRestoreListener stateRestoreListener) { this.topology = topology; this.consumer = consumer; this.stateDirectory = stateDirectory; this.baseDir = stateDirectory.globalStateDir(); this.checkpoint = new OffsetCheckpoint(new File(this.baseDir, CHECKPOINT_FILE_NAME)); + this.stateRestoreListener = stateRestoreListener; } @Override @@ -135,7 +139,7 @@ public void register(final StateStore store, final List topicPartitions = topicPartitionsForStore(store); final Map highWatermarks = consumer.endOffsets(topicPartitions); try { - restoreState(stateRestoreCallback, topicPartitions, highWatermarks); + restoreState(stateRestoreCallback, topicPartitions, highWatermarks, store.name()); stores.put(store.name(), store); } finally { consumer.assign(Collections.emptyList()); @@ -159,7 +163,8 @@ private List topicPartitionsForStore(final StateStore store) { private void restoreState(final StateRestoreCallback stateRestoreCallback, final List topicPartitions, - final Map highWatermarks) { + final Map highWatermarks, + final String storeName) { for (final TopicPartition topicPartition : topicPartitions) { consumer.assign(Collections.singletonList(topicPartition)); final Long checkpoint = checkpointableOffsets.get(topicPartition); @@ -178,6 +183,9 @@ private void restoreState(final StateRestoreCallback stateRestoreCallback, ? stateRestoreCallback : new WrappedBatchingStateRestoreCallback(stateRestoreCallback)); + stateRestoreListener.onRestoreStart(topicPartition, storeName, offset, highWatermark); + long restoreCount = 0L; + while (offset < highWatermark) { final ConsumerRecords records = consumer.poll(100); final List> restoreRecords = new ArrayList<>(); @@ -188,7 +196,10 @@ private void restoreState(final StateRestoreCallback stateRestoreCallback, } } stateRestoreAdapter.restoreAll(restoreRecords); + stateRestoreListener.onBatchRestored(topicPartition, storeName, offset, restoreRecords.size()); + restoreCount += restoreRecords.size(); } + stateRestoreListener.onRestoreEnd(topicPartition, storeName, restoreCount); checkpointableOffsets.put(topicPartition, offset); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStreamThread.java index 41ebcca13da5c..ac742c9ce5f91 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStreamThread.java @@ -27,15 +27,16 @@ import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.StreamsMetrics; import org.apache.kafka.streams.errors.StreamsException; +import org.apache.kafka.streams.processor.StateRestoreListener; import org.apache.kafka.streams.state.internals.ThreadCache; import org.slf4j.Logger; import java.io.IOException; -import java.util.Map; -import java.util.Set; -import java.util.HashSet; import java.util.Arrays; import java.util.Collections; +import java.util.HashSet; +import java.util.Map; +import java.util.Set; import static org.apache.kafka.streams.processor.internals.GlobalStreamThread.State.DEAD; import static org.apache.kafka.streams.processor.internals.GlobalStreamThread.State.PENDING_SHUTDOWN; @@ -113,6 +114,7 @@ public boolean isValidTransition(final ThreadStateTransitionValidator newState) private final Object stateLock = new Object(); private StreamThread.StateListener stateListener = null; private final String logPrefix; + private final StateRestoreListener stateRestoreListener; /** * Set the {@link StreamThread.StateListener} to be notified when state changes. Note this API is internal to @@ -175,7 +177,8 @@ public GlobalStreamThread(final ProcessorTopology topology, final StateDirectory stateDirectory, final Metrics metrics, final Time time, - final String threadClientId) { + final String threadClientId, + final StateRestoreListener stateRestoreListener) { super(threadClientId); this.time = time; this.config = config; @@ -189,6 +192,7 @@ public GlobalStreamThread(final ProcessorTopology topology, this.logContext = new LogContext(logPrefix); this.log = logContext.logger(getClass()); this.cache = new ThreadCache(logContext, cacheSizeBytes, streamsMetrics); + this.stateRestoreListener = stateRestoreListener; } @@ -294,7 +298,7 @@ public void run() { private StateConsumer initialize() { try { - final GlobalStateManager stateMgr = new GlobalStateManagerImpl(topology, consumer, stateDirectory); + final GlobalStateManager stateMgr = new GlobalStateManagerImpl(topology, consumer, stateDirectory, stateRestoreListener); final StateConsumer stateConsumer = new StateConsumer(this.logContext, consumer, diff --git a/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java b/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java index f1ae6dad047b5..baeedf1cf0728 100644 --- a/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java @@ -28,9 +28,9 @@ import org.apache.kafka.streams.integration.utils.IntegrationTestUtils; import org.apache.kafka.streams.kstream.ForeachAction; import org.apache.kafka.streams.processor.StreamPartitioner; +import org.apache.kafka.streams.processor.ThreadMetadata; import org.apache.kafka.streams.processor.internals.GlobalStreamThread; import org.apache.kafka.streams.processor.internals.StreamThread; -import org.apache.kafka.streams.processor.ThreadMetadata; import org.apache.kafka.test.IntegrationTest; import org.apache.kafka.test.MockMetricsReporter; import org.apache.kafka.test.MockStateRestoreListener; @@ -54,9 +54,9 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; -import static org.junit.Assert.assertNotNull; @Category({IntegrationTest.class}) public class KafkaStreamsTest { @@ -179,6 +179,7 @@ public void testStateGlobalThreadClose() throws Exception { builder.globalTable("anyTopic"); props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, numThreads); final KafkaStreams streams = new KafkaStreams(builder.build(), props); + streams.setGlobalStateRestoreListener(new MockStateRestoreListener()); streams.start(); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java index e530d605e227b..2a15ae40b9225 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java @@ -31,6 +31,7 @@ import org.apache.kafka.streams.processor.StateStore; import org.apache.kafka.streams.state.internals.OffsetCheckpoint; import org.apache.kafka.test.MockProcessorNode; +import org.apache.kafka.test.MockStateRestoreListener; import org.apache.kafka.test.NoOpProcessorContext; import org.apache.kafka.test.NoOpReadOnlyStore; import org.apache.kafka.test.TestUtils; @@ -61,6 +62,7 @@ public class GlobalStateManagerImplTest { private final MockTime time = new MockTime(); private final TheStateRestoreCallback stateRestoreCallback = new TheStateRestoreCallback(); + private final MockStateRestoreListener stateRestoreListener = new MockStateRestoreListener(); private final TopicPartition t1 = new TopicPartition("t1", 1); private final TopicPartition t2 = new TopicPartition("t2", 1); private GlobalStateManagerImpl stateManager; @@ -95,7 +97,7 @@ public void before() throws IOException { stateDirPath = TestUtils.tempDirectory().getPath(); stateDirectory = new StateDirectory("appId", stateDirPath, time); consumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST); - stateManager = new GlobalStateManagerImpl(topology, consumer, stateDirectory); + stateManager = new GlobalStateManagerImpl(topology, consumer, stateDirectory, stateRestoreListener); checkpointFile = new File(stateManager.baseDir(), ProcessorStateManager.CHECKPOINT_FILE_NAME); } @@ -452,7 +454,7 @@ public void shouldThrowLockExceptionIfIOExceptionCaughtWhenTryingToLockStateDir( public boolean lockGlobalState(final int retry) throws IOException { throw new IOException("KABOOM!"); } - }); + }, null); try { stateManager.initialize(context); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java index 9d0b6376d5b09..6d9433a93ddae 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java @@ -27,6 +27,7 @@ import org.apache.kafka.streams.errors.StreamsException; import org.apache.kafka.streams.kstream.KStreamBuilder; import org.apache.kafka.streams.processor.StateStore; +import org.apache.kafka.test.MockStateRestoreListener; import org.apache.kafka.test.TestCondition; import org.apache.kafka.test.TestUtils; import org.junit.Before; @@ -65,7 +66,8 @@ public void before() { new StateDirectory("appId", TestUtils.tempDirectory().getPath(), time), new Metrics(), new MockTime(), - "clientId"); + "clientId", + new MockStateRestoreListener()); } @Test @@ -96,7 +98,7 @@ public List partitionsFor(final String topic) { new StateDirectory("appId", TestUtils.tempDirectory().getPath(), time), new Metrics(), new MockTime(), - "clientId"); + "clientId", null); try { globalStreamThread.start(); diff --git a/streams/src/test/java/org/apache/kafka/test/ProcessorTopologyTestDriver.java b/streams/src/test/java/org/apache/kafka/test/ProcessorTopologyTestDriver.java index 27e930997f0b1..f0ca555efeeee 100644 --- a/streams/src/test/java/org/apache/kafka/test/ProcessorTopologyTestDriver.java +++ b/streams/src/test/java/org/apache/kafka/test/ProcessorTopologyTestDriver.java @@ -211,6 +211,7 @@ public List partitionsFor(final String topic) { if (globalTopology != null) { final MockConsumer globalConsumer = createGlobalConsumer(); + final MockStateRestoreListener stateRestoreListener = new MockStateRestoreListener(); for (final String topicName : globalTopology.sourceTopics()) { final List partitionInfos = new ArrayList<>(); partitionInfos.add(new PartitionInfo(topicName, 1, null, null, null)); @@ -220,7 +221,7 @@ public List partitionsFor(final String topic) { globalPartitionsByTopic.put(topicName, partition); offsetsByTopicPartition.put(partition, new AtomicLong()); } - final GlobalStateManagerImpl stateManager = new GlobalStateManagerImpl(globalTopology, globalConsumer, stateDirectory); + final GlobalStateManagerImpl stateManager = new GlobalStateManagerImpl(globalTopology, globalConsumer, stateDirectory, stateRestoreListener); globalStateTask = new GlobalStateUpdateTask(globalTopology, new GlobalProcessorContextImpl(config, stateManager, streamsMetrics, cache), stateManager, new LogAndContinueExceptionHandler() From 62c601a97517db24a70c884df4e67e99724af3ec Mon Sep 17 00:00:00 2001 From: Bill Bejeck Date: Wed, 27 Sep 2017 10:01:33 -0400 Subject: [PATCH 2/4] KAFKA-5958: add test --- .../internals/GlobalStateManagerImplTest.java | 21 +++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java index 2a15ae40b9225..fcf79a8b22486 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java @@ -50,6 +50,9 @@ import java.util.Map; import java.util.Set; +import static org.apache.kafka.test.MockStateRestoreListener.RESTORE_BATCH; +import static org.apache.kafka.test.MockStateRestoreListener.RESTORE_END; +import static org.apache.kafka.test.MockStateRestoreListener.RESTORE_START; import static org.hamcrest.CoreMatchers.equalTo; import static org.hamcrest.MatcherAssert.assertThat; import static org.junit.Assert.assertEquals; @@ -205,6 +208,24 @@ public void shouldRestoreRecordsUpToHighwatermark() { assertEquals(2, stateRestoreCallback.restored.size()); } + @Test + public void shouldListenForRestoreEvents() { + initializeConsumer(5, 1, t1); + stateManager.initialize(context); + + final TheStateRestoreCallback stateRestoreCallback = new TheStateRestoreCallback(); + stateManager.register(store1, false, stateRestoreCallback); + + assertThat(stateRestoreListener.restoreStartOffset, equalTo(1L)); + assertThat(stateRestoreListener.restoreEndOffset, equalTo(5L)); + assertThat(stateRestoreListener.totalNumRestored, equalTo(5L)); + + + assertThat(stateRestoreListener.storeNameCalledStates.get(RESTORE_START), equalTo(store1.name())); + assertThat(stateRestoreListener.storeNameCalledStates.get(RESTORE_BATCH), equalTo(store1.name())); + assertThat(stateRestoreListener.storeNameCalledStates.get(RESTORE_END), equalTo(store1.name())); + } + @Test public void shouldRestoreRecordsFromCheckpointToHighwatermark() throws IOException { initializeConsumer(5, 6, t1); From 5430884fecf1e34690042f207d4f3884affdba83 Mon Sep 17 00:00:00 2001 From: Bill Bejeck Date: Wed, 27 Sep 2017 11:37:12 -0400 Subject: [PATCH 3/4] KAFKA-5958: updates for comments --- .../streams/processor/internals/GlobalStreamThread.java | 5 ++++- .../java/org/apache/kafka/streams/KafkaStreamsTest.java | 2 +- .../processor/internals/GlobalStateManagerImplTest.java | 2 +- .../streams/processor/internals/GlobalStreamThreadTest.java | 6 ++++-- .../org/apache/kafka/test/ProcessorTopologyTestDriver.java | 5 ++++- 5 files changed, 14 insertions(+), 6 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStreamThread.java index ac742c9ce5f91..a365addad0f36 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/GlobalStreamThread.java @@ -298,7 +298,10 @@ public void run() { private StateConsumer initialize() { try { - final GlobalStateManager stateMgr = new GlobalStateManagerImpl(topology, consumer, stateDirectory, stateRestoreListener); + final GlobalStateManager stateMgr = new GlobalStateManagerImpl(topology, + consumer, + stateDirectory, + stateRestoreListener); final StateConsumer stateConsumer = new StateConsumer(this.logContext, consumer, diff --git a/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java b/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java index baeedf1cf0728..6625976d9ffe9 100644 --- a/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java @@ -179,7 +179,7 @@ public void testStateGlobalThreadClose() throws Exception { builder.globalTable("anyTopic"); props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, numThreads); final KafkaStreams streams = new KafkaStreams(builder.build(), props); - streams.setGlobalStateRestoreListener(new MockStateRestoreListener()); + streams.setGlobalStateRestoreListener(null); streams.start(); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java index fcf79a8b22486..b438347c45d25 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStateManagerImplTest.java @@ -475,7 +475,7 @@ public void shouldThrowLockExceptionIfIOExceptionCaughtWhenTryingToLockStateDir( public boolean lockGlobalState(final int retry) throws IOException { throw new IOException("KABOOM!"); } - }, null); + }, stateRestoreListener); try { stateManager.initialize(context); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java index 6d9433a93ddae..29f1ac08fba37 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/GlobalStreamThreadTest.java @@ -50,6 +50,7 @@ public class GlobalStreamThreadTest { private final KStreamBuilder builder = new KStreamBuilder(); private final MockConsumer mockConsumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST); private final MockTime time = new MockTime(); + private final MockStateRestoreListener stateRestoreListener = new MockStateRestoreListener(); private GlobalStreamThread globalStreamThread; private StreamsConfig config; @@ -67,7 +68,7 @@ public void before() { new Metrics(), new MockTime(), "clientId", - new MockStateRestoreListener()); + stateRestoreListener); } @Test @@ -98,7 +99,8 @@ public List partitionsFor(final String topic) { new StateDirectory("appId", TestUtils.tempDirectory().getPath(), time), new Metrics(), new MockTime(), - "clientId", null); + "clientId", + stateRestoreListener); try { globalStreamThread.start(); diff --git a/streams/src/test/java/org/apache/kafka/test/ProcessorTopologyTestDriver.java b/streams/src/test/java/org/apache/kafka/test/ProcessorTopologyTestDriver.java index f0ca555efeeee..babf704ba8769 100644 --- a/streams/src/test/java/org/apache/kafka/test/ProcessorTopologyTestDriver.java +++ b/streams/src/test/java/org/apache/kafka/test/ProcessorTopologyTestDriver.java @@ -221,7 +221,10 @@ public List partitionsFor(final String topic) { globalPartitionsByTopic.put(topicName, partition); offsetsByTopicPartition.put(partition, new AtomicLong()); } - final GlobalStateManagerImpl stateManager = new GlobalStateManagerImpl(globalTopology, globalConsumer, stateDirectory, stateRestoreListener); + final GlobalStateManagerImpl stateManager = new GlobalStateManagerImpl(globalTopology, + globalConsumer, + stateDirectory, + stateRestoreListener); globalStateTask = new GlobalStateUpdateTask(globalTopology, new GlobalProcessorContextImpl(config, stateManager, streamsMetrics, cache), stateManager, new LogAndContinueExceptionHandler() From f9e08be606d89bbe8f11677e0f6a528a88bdc756 Mon Sep 17 00:00:00 2001 From: Bill Bejeck Date: Wed, 27 Sep 2017 12:04:10 -0400 Subject: [PATCH 4/4] KAFKA-5958: remove unneeded setter --- .../src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java b/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java index 6625976d9ffe9..4bd289085adf0 100644 --- a/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java @@ -179,7 +179,6 @@ public void testStateGlobalThreadClose() throws Exception { builder.globalTable("anyTopic"); props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, numThreads); final KafkaStreams streams = new KafkaStreams(builder.build(), props); - streams.setGlobalStateRestoreListener(null); streams.start();