diff --git a/checkstyle/suppressions.xml b/checkstyle/suppressions.xml index 1bb797bd4497c..f3c463457c01e 100644 --- a/checkstyle/suppressions.xml +++ b/checkstyle/suppressions.xml @@ -181,7 +181,7 @@ files="StreamsPartitionAssignor.java"/> + files="(AssignorConfiguration|InternalTopologyBuilder|KafkaStreams|ProcessorStateManager|StreamsPartitionAssignor|StreamThread|TaskManager).java"/> 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 34a2d4d322067..f46e4013b5e67 100644 --- a/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java +++ b/streams/src/main/java/org/apache/kafka/streams/KafkaStreams.java @@ -38,6 +38,7 @@ import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.errors.ProcessorStateException; import org.apache.kafka.streams.errors.StreamsException; +import org.apache.kafka.streams.errors.TopologyException; import org.apache.kafka.streams.internals.ApiUtils; import org.apache.kafka.streams.internals.metrics.ClientMetrics; import org.apache.kafka.streams.kstream.KStream; @@ -51,6 +52,7 @@ import org.apache.kafka.streams.processor.internals.ClientUtils; import org.apache.kafka.streams.processor.internals.DefaultKafkaClientSupplier; import org.apache.kafka.streams.processor.internals.GlobalStreamThread; +import org.apache.kafka.streams.processor.internals.GlobalStreamThread.State; import org.apache.kafka.streams.processor.internals.InternalTopologyBuilder; import org.apache.kafka.streams.processor.internals.ProcessorTopology; import org.apache.kafka.streams.processor.internals.StateDirectory; @@ -482,8 +484,9 @@ public synchronized void onChange(final Thread thread, final GlobalStreamThread.State newState = (GlobalStreamThread.State) abstractNewState; globalThreadState = newState; - // special case when global thread is dead - if (newState == GlobalStreamThread.State.DEAD) { + if (newState == GlobalStreamThread.State.RUNNING) { + maybeSetRunning(); + } else if (newState == GlobalStreamThread.State.DEAD) { if (setState(State.ERROR)) { log.error("Global thread has died. The instance will be in error state and should be closed."); } @@ -701,18 +704,34 @@ private KafkaStreams(final InternalTopologyBuilder internalTopologyBuilder, internalTopologyBuilder, parseHostInfo(config.getString(StreamsConfig.APPLICATION_SERVER_CONFIG))); + final int numStreamThreads; + if (internalTopologyBuilder.hasNoNonGlobalTopology()) { + log.info("Overriding number of StreamThreads to zero for global-only topology"); + numStreamThreads = 0; + } else { + numStreamThreads = config.getInt(StreamsConfig.NUM_STREAM_THREADS_CONFIG); + } + // create the stream thread, global update thread, and cleanup thread - threads = new StreamThread[config.getInt(StreamsConfig.NUM_STREAM_THREADS_CONFIG)]; + threads = new StreamThread[numStreamThreads]; + + final ProcessorTopology globalTaskTopology = internalTopologyBuilder.buildGlobalStateTopology(); + final boolean hasGlobalTopology = globalTaskTopology != null; + + if (numStreamThreads == 0 && !hasGlobalTopology) { + log.error("Topology with no input topics will create no stream threads and no global thread."); + throw new TopologyException("Topology has no stream threads and no global threads, " + + "must subscribe to at least one source topic or global table."); + } long totalCacheSize = config.getLong(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG); if (totalCacheSize < 0) { totalCacheSize = 0; log.warn("Negative cache size passed in. Reverting to cache size of 0 bytes."); } - final ProcessorTopology globalTaskTopology = internalTopologyBuilder.buildGlobalStateTopology(); - final long cacheSizePerThread = totalCacheSize / (threads.length + (globalTaskTopology == null ? 0 : 1)); + final long cacheSizePerThread = totalCacheSize / (threads.length + (hasGlobalTopology ? 1 : 0)); final boolean hasPersistentStores = taskTopology.hasPersistentLocalStore() || - (globalTaskTopology != null && globalTaskTopology.hasPersistentGlobalStore()); + (hasGlobalTopology && globalTaskTopology.hasPersistentGlobalStore()); try { stateDirectory = new StateDirectory(config, time, hasPersistentStores); @@ -722,7 +741,7 @@ private KafkaStreams(final InternalTopologyBuilder internalTopologyBuilder, final StateRestoreListener delegatingStateRestoreListener = new DelegatingStateRestoreListener(); GlobalStreamThread.State globalThreadState = null; - if (globalTaskTopology != null) { + if (hasGlobalTopology) { final String globalThreadId = clientId + "-GlobalStreamThread"; globalStreamThread = new GlobalStreamThread( globalTaskTopology, @@ -766,7 +785,7 @@ private KafkaStreams(final InternalTopologyBuilder internalTopologyBuilder, Math.toIntExact(Arrays.stream(threads).filter(thread -> thread.state().isAlive()).count())); final StreamStateListener streamStateListener = new StreamStateListener(threadState, globalThreadState); - if (globalTaskTopology != null) { + if (hasGlobalTopology) { globalStreamThread.setStateListener(streamStateListener); } for (final StreamThread thread : threads) { diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopologyBuilder.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopologyBuilder.java index 854ce1cb81714..271951796daed 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopologyBuilder.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/InternalTopologyBuilder.java @@ -1259,6 +1259,10 @@ synchronized Pattern sourceTopicPattern() { return sourceTopicPattern; } + public boolean hasNoNonGlobalTopology() { + return !usesPatternSubscription() && sourceTopicCollection().isEmpty(); + } + private boolean isGlobalSource(final String nodeName) { final NodeFactory nodeFactory = nodeFactories.get(nodeName); 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 49fa105908609..05271e82b07a9 100644 --- a/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/KafkaStreamsTest.java @@ -33,6 +33,7 @@ import org.apache.kafka.common.serialization.StringSerializer; import org.apache.kafka.common.utils.MockTime; import org.apache.kafka.common.utils.Time; +import org.apache.kafka.streams.errors.TopologyException; import org.apache.kafka.streams.internals.metrics.ClientMetrics; import org.apache.kafka.streams.kstream.Materialized; import org.apache.kafka.streams.processor.AbstractProcessor; @@ -90,6 +91,7 @@ import static org.easymock.EasyMock.capture; import static org.hamcrest.CoreMatchers.hasItem; import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.equalTo; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNull; @@ -111,7 +113,7 @@ public class KafkaStreamsTest { private MockTime time; private Properties props; - + @Mock private StateDirectory stateDirectory; @Mock @@ -247,6 +249,7 @@ private void prepareStreams() throws Exception { globalStreamThread.shutdown(); EasyMock.expectLastCall().andAnswer(() -> { supplier.restoreConsumer.close(); + for (final MockProducer producer : supplier.producers) { producer.close(); } @@ -327,7 +330,7 @@ private void prepareStreamThread(final StreamThread thread, final boolean termin @Test public void testShouldTransitToNotRunningIfCloseRightAfterCreated() { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); streams.close(); Assert.assertEquals(KafkaStreams.State.NOT_RUNNING, streams.state()); @@ -335,7 +338,7 @@ public void testShouldTransitToNotRunningIfCloseRightAfterCreated() { @Test public void stateShouldTransitToRunningIfNonDeadThreadsBackToRunning() throws InterruptedException { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); streams.setStateListener(streamsStateListener); Assert.assertEquals(0, streamsStateListener.numChanges); @@ -403,7 +406,7 @@ public void stateShouldTransitToRunningIfNonDeadThreadsBackToRunning() throws In @Test public void stateShouldTransitToErrorIfAllThreadsDead() throws InterruptedException { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); streams.setStateListener(streamsStateListener); Assert.assertEquals(0, streamsStateListener.numChanges); @@ -467,7 +470,7 @@ public void stateShouldTransitToErrorIfAllThreadsDead() throws InterruptedExcept @Test public void shouldCleanupResourcesOnCloseWithoutPreviousStart() throws Exception { - final StreamsBuilder builder = new StreamsBuilder(); + final StreamsBuilder builder = getBuilderWithSource(); builder.globalTable("anyTopic"); final KafkaStreams streams = new KafkaStreams(builder.build(), props, supplier, time); @@ -487,7 +490,7 @@ public void shouldCleanupResourcesOnCloseWithoutPreviousStart() throws Exception @Test public void testStateThreadClose() throws Exception { // make sure we have the global state thread running too - final StreamsBuilder builder = new StreamsBuilder(); + final StreamsBuilder builder = getBuilderWithSource(); builder.globalTable("anyTopic"); final KafkaStreams streams = new KafkaStreams(builder.build(), props, supplier, time); @@ -524,7 +527,7 @@ public void testStateThreadClose() throws Exception { @Test public void testStateGlobalThreadClose() throws Exception { // make sure we have the global state thread running too - final StreamsBuilder builder = new StreamsBuilder(); + final StreamsBuilder builder = getBuilderWithSource(); builder.globalTable("anyTopic"); final KafkaStreams streams = new KafkaStreams(builder.build(), props, supplier, time); @@ -552,7 +555,7 @@ public void testStateGlobalThreadClose() throws Exception { public void testInitializesAndDestroysMetricsReporters() { final int oldInitCount = MockMetricsReporter.INIT_COUNT.get(); - try (final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time)) { + try (final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time)) { final int newInitCount = MockMetricsReporter.INIT_COUNT.get(); final int initDiff = newInitCount - oldInitCount; assertTrue("some reporters should be initialized by calling on construction", initDiff > 0); @@ -566,7 +569,7 @@ public void testInitializesAndDestroysMetricsReporters() { @Test public void testCloseIsIdempotent() { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); streams.close(); final int closeCount = MockMetricsReporter.CLOSE_COUNT.get(); @@ -577,7 +580,7 @@ public void testCloseIsIdempotent() { @Test public void testCannotStartOnceClosed() { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); streams.start(); streams.close(); try { @@ -592,7 +595,7 @@ public void testCannotStartOnceClosed() { @Test public void shouldNotSetGlobalRestoreListenerAfterStarting() { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); streams.start(); try { streams.setGlobalStateRestoreListener(null); @@ -606,7 +609,7 @@ public void shouldNotSetGlobalRestoreListenerAfterStarting() { @Test public void shouldThrowExceptionSettingUncaughtExceptionHandlerNotInCreateState() { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); streams.start(); try { streams.setUncaughtExceptionHandler(null); @@ -618,7 +621,7 @@ public void shouldThrowExceptionSettingUncaughtExceptionHandlerNotInCreateState( @Test public void shouldThrowExceptionSettingStateListenerNotInCreateState() { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); streams.start(); try { streams.setStateListener(null); @@ -630,7 +633,7 @@ public void shouldThrowExceptionSettingStateListenerNotInCreateState() { @Test public void shouldAllowCleanupBeforeStartAndAfterClose() { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); try { streams.cleanUp(); streams.start(); @@ -642,7 +645,7 @@ public void shouldAllowCleanupBeforeStartAndAfterClose() { @Test public void shouldThrowOnCleanupWhileRunning() throws InterruptedException { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); streams.start(); TestUtils.waitForCondition( () -> streams.state() == KafkaStreams.State.RUNNING, @@ -658,32 +661,32 @@ public void shouldThrowOnCleanupWhileRunning() throws InterruptedException { @Test(expected = IllegalStateException.class) public void shouldNotGetAllTasksWhenNotRunning() { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); streams.allMetadata(); } @Test(expected = IllegalStateException.class) public void shouldNotGetAllTasksWithStoreWhenNotRunning() { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); streams.allMetadataForStore("store"); } @Test(expected = IllegalStateException.class) public void shouldNotGetQueryMetadataWithSerializerWhenNotRunningOrRebalancing() { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); streams.queryMetadataForKey("store", "key", Serdes.String().serializer()); } @Test public void shouldGetQueryMetadataWithSerializerWhenRunningOrRebalancing() { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); streams.start(); assertEquals(KeyQueryMetadata.NOT_AVAILABLE, streams.queryMetadataForKey("store", "key", Serdes.String().serializer())); } @Test(expected = IllegalStateException.class) public void shouldNotGetQueryMetadataWithPartitionerWhenNotRunningOrRebalancing() { - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); streams.queryMetadataForKey("store", "key", (topic, key, value, numPartitions) -> 0); } @@ -703,7 +706,7 @@ public void shouldReturnEmptyLocalStorePartitionLags() { EasyMock.expect(mockClientSupplier.getAdmin(anyObject())).andReturn(mockAdminClient); EasyMock.replay(result, mockAdminClient, mockClientSupplier); - final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, mockClientSupplier, time); + final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, mockClientSupplier, time); streams.start(); assertEquals(0, streams.allLocalStorePartitionLags().size()); } @@ -711,21 +714,21 @@ public void shouldReturnEmptyLocalStorePartitionLags() { @Test public void shouldReturnFalseOnCloseWhenThreadsHaventTerminated() { // do not use mock time so that it can really elapse - try (final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier)) { + try (final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier)) { assertFalse(streams.close(Duration.ofMillis(10L))); } } @Test(expected = IllegalArgumentException.class) public void shouldThrowOnNegativeTimeoutForClose() { - try (final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time)) { + try (final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time)) { streams.close(Duration.ofMillis(-1L)); } } @Test public void shouldNotBlockInCloseForZeroDuration() { - try (final KafkaStreams streams = new KafkaStreams(new StreamsBuilder().build(), props, supplier, time)) { + try (final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time)) { // with mock time that does not elapse, close would not return if it ever waits on the state transition assertFalse(streams.close(Duration.ZERO)); } @@ -758,7 +761,7 @@ public void shouldTriggerRecordingOfRocksDBMetricsIfRecordingLevelIsDebug() { builder.table("topic", Materialized.as("store")); props.setProperty(StreamsConfig.METRICS_RECORDING_LEVEL_CONFIG, RecordingLevel.DEBUG.name()); - try (final KafkaStreams streams = new KafkaStreams(builder.build(), props, supplier, time)) { + try (final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time)) { streams.start(); } @@ -781,7 +784,7 @@ public void shouldNotTriggerRecordingOfRocksDBMetricsIfRecordingLevelIsInfo() { builder.table("topic", Materialized.as("store")); props.setProperty(StreamsConfig.METRICS_RECORDING_LEVEL_CONFIG, RecordingLevel.INFO.name()); - try (final KafkaStreams streams = new KafkaStreams(builder.build(), props, supplier, time)) { + try (final KafkaStreams streams = new KafkaStreams(getBuilderWithSource().build(), props, supplier, time)) { streams.start(); } @@ -801,7 +804,7 @@ public void shouldWarnAboutRocksDBConfigSetterIsNotGuaranteedToBeBackwardsCompat props.setProperty(StreamsConfig.ROCKSDB_CONFIG_SETTER_CLASS_CONFIG, TestRocksDbConfigSetter.class.getName()); try (final LogCaptureAppender appender = LogCaptureAppender.createAndRegister()) { - new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + new KafkaStreams(getBuilderWithSource().build(), props, supplier, time); assertThat(appender.getMessages(), hasItem("stream-client [" + CLIENT_ID + "] " + "RocksDB's version will be bumped to version 6+ via KAFKA-8897 in a future release. " @@ -883,6 +886,49 @@ public void statefulTopologyShouldCreateStateDirectory() throws Exception { startStreamsAndCheckDirExists(topology, true); } + @Test + public void shouldThrowTopologyExceptionOnEmptyTopology() { + try { + new KafkaStreams(new StreamsBuilder().build(), props, supplier, time); + fail("Should have thrown TopologyException"); + } catch (final TopologyException e) { + assertThat( + e.getMessage(), + equalTo("Invalid topology: Topology has no stream threads and no global threads, " + + "must subscribe to at least one source topic or global table.")); + } + } + + @Test + public void shouldNotCreateStreamThreadsForGlobalOnlyTopology() { + final StreamsBuilder builder = new StreamsBuilder(); + builder.globalTable("anyTopic"); + final KafkaStreams streams = new KafkaStreams(builder.build(), props, supplier, time); + + assertThat(streams.threads.length, equalTo(0)); + } + + @Test + public void shouldTransitToRunningWithGlobalOnlyTopology() throws InterruptedException { + final StreamsBuilder builder = new StreamsBuilder(); + builder.globalTable("anyTopic"); + final KafkaStreams streams = new KafkaStreams(builder.build(), props, supplier, time); + + assertThat(streams.threads.length, equalTo(0)); + assertEquals(streams.state(), KafkaStreams.State.CREATED); + + streams.start(); + TestUtils.waitForCondition( + () -> streams.state() == KafkaStreams.State.RUNNING, + "Streams never started, state is " + streams.state()); + + streams.close(); + + TestUtils.waitForCondition( + () -> streams.state() == KafkaStreams.State.NOT_RUNNING, + "Streams never stopped."); + } + @SuppressWarnings("unchecked") private Topology getStatefulTopology(final String inputTopic, final String outputTopic, @@ -925,6 +971,12 @@ public void process(final String key, final String value) { new MockProcessorSupplier<>()); return topology; } + + private StreamsBuilder getBuilderWithSource() { + final StreamsBuilder builder = new StreamsBuilder(); + builder.stream("source-topic"); + return builder; + } private void startStreamsAndCheckDirExists(final Topology topology, final boolean shouldFilesExist) throws Exception { diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/GlobalKTableIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/GlobalKTableIntegrationTest.java index 628d913747831..88ef238c7f074 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/GlobalKTableIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/GlobalKTableIntegrationTest.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.streams.integration; +import java.time.Duration; import kafka.utils.MockTime; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.common.serialization.LongSerializer; @@ -23,6 +24,7 @@ import org.apache.kafka.common.serialization.StringSerializer; import org.apache.kafka.common.utils.Bytes; import org.apache.kafka.streams.KafkaStreams; +import org.apache.kafka.streams.KafkaStreams.State; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; @@ -56,6 +58,8 @@ import java.util.Map; import java.util.Properties; +import static java.util.Collections.singletonList; +import static org.apache.kafka.streams.integration.utils.IntegrationTestUtils.waitForApplicationState; import static org.hamcrest.MatcherAssert.assertThat; import static org.hamcrest.core.IsEqual.equalTo; @@ -271,6 +275,20 @@ public void shouldRestoreGlobalInMemoryKTableOnRestart() throws Exception { assertThat(timestampedStore.approximateNumEntries(), equalTo(4L)); } + @Test + public void shouldGetToRunningWithOnlyGlobalTopology() throws Exception { + builder = new StreamsBuilder(); + globalTable = builder.globalTable( + globalTableTopic, + Consumed.with(Serdes.Long(), Serdes.String()), + Materialized.as(Stores.inMemoryKeyValueStore(globalStore))); + + startStreams(); + waitForApplicationState(singletonList(kafkaStreams), State.RUNNING, Duration.ofSeconds(30)); + + kafkaStreams.close(); + } + private void createTopics() throws Exception { streamTopic = "stream-" + testName.getMethodName(); globalTableTopic = "globalTable-" + testName.getMethodName();