diff --git a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java index 2ec1d3dc51d19..eadc3dff4e5cd 100644 --- a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java +++ b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java @@ -872,6 +872,9 @@ public static class InternalConfig { public static final String ASSIGNMENT_ERROR_CODE = "__assignment.error.code__"; public static final String NEXT_SCHEDULED_REBALANCE_MS = "__next.probing.rebalance.ms__"; public static final String TIME = "__time__"; + + // This is settable in the main Streams config, but it's a private API for testing + public static final String ASSIGNMENT_LISTENER = "__asignment.listener__"; } /** diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java index 05259972c961a..1441612d40f60 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignor.java @@ -36,6 +36,7 @@ import org.apache.kafka.streams.processor.internals.assignment.AssignmentInfo; import org.apache.kafka.streams.processor.internals.assignment.AssignorConfiguration; import org.apache.kafka.streams.processor.internals.assignment.AssignorConfiguration.AssignmentConfigs; +import org.apache.kafka.streams.processor.internals.assignment.AssignorConfiguration.AssignmentListener; import org.apache.kafka.streams.processor.internals.assignment.AssignorError; import org.apache.kafka.streams.processor.internals.assignment.ClientState; import org.apache.kafka.streams.processor.internals.assignment.CopartitionedTopicsEnforcer; @@ -171,6 +172,7 @@ public String toString() { private InternalTopicManager internalTopicManager; private CopartitionedTopicsEnforcer copartitionedTopicsEnforcer; private RebalanceProtocol rebalanceProtocol; + private AssignmentListener assignmentListener; private Supplier taskAssignorSupplier; @@ -189,20 +191,21 @@ public void configure(final Map configs) { log = new LogContext(logPrefix).logger(getClass()); usedSubscriptionMetadataVersion = assignorConfiguration .configuredMetadataVersion(usedSubscriptionMetadataVersion); - taskManager = assignorConfiguration.getTaskManager(); - streamsMetadataState = assignorConfiguration.getStreamsMetadataState(); - assignmentErrorCode = assignorConfiguration.getAssignmentErrorCode(configs); - nextScheduledRebalanceMs = assignorConfiguration.getNextScheduledRebalanceMs(configs); - time = assignorConfiguration.getTime(configs); - assignmentConfigs = assignorConfiguration.getAssignmentConfigs(); - partitionGrouper = assignorConfiguration.getPartitionGrouper(); - userEndPoint = assignorConfiguration.getUserEndPoint(); - adminClient = assignorConfiguration.getAdminClient(); - adminClientTimeout = assignorConfiguration.getAdminClientTimeout(); - internalTopicManager = assignorConfiguration.getInternalTopicManager(); - copartitionedTopicsEnforcer = assignorConfiguration.getCopartitionedTopicsEnforcer(); + taskManager = assignorConfiguration.taskManager(); + streamsMetadataState = assignorConfiguration.streamsMetadataState(); + assignmentErrorCode = assignorConfiguration.assignmentErrorCode(); + nextScheduledRebalanceMs = assignorConfiguration.nextScheduledRebalanceMs(); + time = assignorConfiguration.time(); + assignmentConfigs = assignorConfiguration.assignmentConfigs(); + partitionGrouper = assignorConfiguration.partitionGrouper(); + userEndPoint = assignorConfiguration.userEndPoint(); + adminClient = assignorConfiguration.adminClient(); + adminClientTimeout = assignorConfiguration.adminClientTimeout(); + internalTopicManager = assignorConfiguration.internalTopicManager(); + copartitionedTopicsEnforcer = assignorConfiguration.copartitionedTopicsEnforcer(); rebalanceProtocol = assignorConfiguration.rebalanceProtocol(); - taskAssignorSupplier = assignorConfiguration::getTaskAssignor; + taskAssignorSupplier = assignorConfiguration::taskAssignor; + assignmentListener = assignorConfiguration.assignmentListener(); } @Override @@ -913,8 +916,10 @@ private Map computeNewAssignment(final Map internalConfigs; - @SuppressWarnings("deprecation") public AssignorConfiguration(final Map configs) { streamsConfig = new QuietStreamsConfig(configs); + internalConfigs = configs; // Setting the logger with the passed in client thread name logPrefix = String.format("stream-thread [%s] ", streamsConfig.getString(CommonClientConfigs.CLIENT_ID_CONFIG)); final LogContext logContext = new LogContext(logPrefix); log = logContext.logger(getClass()); - assignmentConfigs = new AssignmentConfigs(streamsConfig); - - partitionGrouper = streamsConfig.getConfiguredInstance( - StreamsConfig.PARTITION_GROUPER_CLASS_CONFIG, - org.apache.kafka.streams.processor.PartitionGrouper.class - ); - - final String configuredUserEndpoint = streamsConfig.getString(StreamsConfig.APPLICATION_SERVER_CONFIG); - if (configuredUserEndpoint != null && !configuredUserEndpoint.isEmpty()) { - try { - final String host = getHost(configuredUserEndpoint); - final Integer port = getPort(configuredUserEndpoint); - - if (host == null || port == null) { - throw new ConfigException( - String.format( - "%s Config %s isn't in the correct format. Expected a host:port pair but received %s", - logPrefix, StreamsConfig.APPLICATION_SERVER_CONFIG, configuredUserEndpoint - ) - ); - } - } catch (final NumberFormatException nfe) { - throw new ConfigException( - String.format("%s Invalid port supplied in %s for config %s: %s", - logPrefix, configuredUserEndpoint, StreamsConfig.APPLICATION_SERVER_CONFIG, nfe) - ); - } - userEndPoint = configuredUserEndpoint; - } else { - userEndPoint = null; - } - { final Object o = configs.get(StreamsConfig.InternalConfig.TASK_MANAGER_FOR_PARTITION_ASSIGNOR); if (o == null) { @@ -119,25 +81,6 @@ public AssignorConfiguration(final Map configs) { taskManager = (TaskManager) o; } - { - final Object o = configs.get(StreamsConfig.InternalConfig.STREAMS_METADATA_STATE_FOR_PARTITION_ASSIGNOR); - if (o == null) { - final KafkaException fatalException = new KafkaException("StreamsMetadataState is not specified"); - log.error(fatalException.getMessage(), fatalException); - throw fatalException; - } - - if (!(o instanceof StreamsMetadataState)) { - final KafkaException fatalException = new KafkaException( - String.format("%s is not an instance of %s", o.getClass().getName(), StreamsMetadataState.class.getName()) - ); - log.error(fatalException.getMessage(), fatalException); - throw fatalException; - } - - streamsMetadataState = (StreamsMetadataState) o; - } - { final Object o = configs.get(StreamsConfig.InternalConfig.STREAMS_ADMIN_CLIENT); if (o == null) { @@ -155,13 +98,8 @@ public AssignorConfiguration(final Map configs) { } adminClient = (Admin) o; - internalTopicManager = new InternalTopicManager(adminClient, streamsConfig); } - adminClientTimeout = streamsConfig.getInt(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG); - - copartitionedTopicsEnforcer = new CopartitionedTopicsEnforcer(logPrefix); - { final String o = (String) configs.get(INTERNAL_TASK_ASSIGNOR_CLASS); if (o == null) { @@ -172,8 +110,8 @@ public AssignorConfiguration(final Map configs) { } } - public AtomicInteger getAssignmentErrorCode(final Map configs) { - final Object ai = configs.get(StreamsConfig.InternalConfig.ASSIGNMENT_ERROR_CODE); + public AtomicInteger assignmentErrorCode() { + final Object ai = internalConfigs.get(StreamsConfig.InternalConfig.ASSIGNMENT_ERROR_CODE); if (ai == null) { final KafkaException fatalException = new KafkaException("assignmentErrorCode is not specified"); log.error(fatalException.getMessage(), fatalException); @@ -190,8 +128,8 @@ public AtomicInteger getAssignmentErrorCode(final Map configs) { return (AtomicInteger) ai; } - public AtomicLong getNextScheduledRebalanceMs(final Map configs) { - final Object al = configs.get(InternalConfig.NEXT_SCHEDULED_REBALANCE_MS); + public AtomicLong nextScheduledRebalanceMs() { + final Object al = internalConfigs.get(InternalConfig.NEXT_SCHEDULED_REBALANCE_MS); if (al == null) { final KafkaException fatalException = new KafkaException("nextProbingRebalanceMs is not specified"); log.error(fatalException.getMessage(), fatalException); @@ -209,8 +147,8 @@ public AtomicLong getNextScheduledRebalanceMs(final Map configs) { return (AtomicLong) al; } - public Time getTime(final Map configs) { - final Object t = configs.get(InternalConfig.TIME); + public Time time() { + final Object t = internalConfigs.get(InternalConfig.TIME); if (t == null) { final KafkaException fatalException = new KafkaException("time is not specified"); log.error(fatalException.getMessage(), fatalException); @@ -228,12 +166,27 @@ public Time getTime(final Map configs) { return (Time) t; } - public TaskManager getTaskManager() { + public TaskManager taskManager() { return taskManager; } - public StreamsMetadataState getStreamsMetadataState() { - return streamsMetadataState; + public StreamsMetadataState streamsMetadataState() { + final Object o = internalConfigs.get(StreamsConfig.InternalConfig.STREAMS_METADATA_STATE_FOR_PARTITION_ASSIGNOR); + if (o == null) { + final KafkaException fatalException = new KafkaException("StreamsMetadataState is not specified"); + log.error(fatalException.getMessage(), fatalException); + throw fatalException; + } + + if (!(o instanceof StreamsMetadataState)) { + final KafkaException fatalException = new KafkaException( + String.format("%s is not an instance of %s", o.getClass().getName(), StreamsMetadataState.class.getName()) + ); + log.error(fatalException.getMessage(), fatalException); + throw fatalException; + } + + return (StreamsMetadataState) o; } public RebalanceProtocol rebalanceProtocol() { @@ -301,35 +254,61 @@ public int configuredMetadataVersion(final int priorVersion) { } @SuppressWarnings("deprecation") - public org.apache.kafka.streams.processor.PartitionGrouper getPartitionGrouper() { - return partitionGrouper; + public org.apache.kafka.streams.processor.PartitionGrouper partitionGrouper() { + return streamsConfig.getConfiguredInstance( + StreamsConfig.PARTITION_GROUPER_CLASS_CONFIG, + org.apache.kafka.streams.processor.PartitionGrouper.class + ); } - public String getUserEndPoint() { - return userEndPoint; + public String userEndPoint() { + final String configuredUserEndpoint = streamsConfig.getString(StreamsConfig.APPLICATION_SERVER_CONFIG); + if (configuredUserEndpoint != null && !configuredUserEndpoint.isEmpty()) { + try { + final String host = getHost(configuredUserEndpoint); + final Integer port = getPort(configuredUserEndpoint); + + if (host == null || port == null) { + throw new ConfigException( + String.format( + "%s Config %s isn't in the correct format. Expected a host:port pair but received %s", + logPrefix, StreamsConfig.APPLICATION_SERVER_CONFIG, configuredUserEndpoint + ) + ); + } + } catch (final NumberFormatException nfe) { + throw new ConfigException( + String.format("%s Invalid port supplied in %s for config %s: %s", + logPrefix, configuredUserEndpoint, StreamsConfig.APPLICATION_SERVER_CONFIG, nfe) + ); + } + return configuredUserEndpoint; + } else { + return null; + } } - public Admin getAdminClient() { + public Admin adminClient() { return adminClient; } - public int getAdminClientTimeout() { - return adminClientTimeout; + public int adminClientTimeout() { + return streamsConfig.getInt(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG); } - public InternalTopicManager getInternalTopicManager() { - return internalTopicManager; + public InternalTopicManager internalTopicManager() { + return new InternalTopicManager(adminClient, streamsConfig); } - public CopartitionedTopicsEnforcer getCopartitionedTopicsEnforcer() { - return copartitionedTopicsEnforcer; + public CopartitionedTopicsEnforcer copartitionedTopicsEnforcer() { + return new CopartitionedTopicsEnforcer(logPrefix); } - public AssignmentConfigs getAssignmentConfigs() { - return assignmentConfigs; + public AssignmentConfigs assignmentConfigs() { + return new AssignmentConfigs(streamsConfig); } - public TaskAssignor getTaskAssignor() { + public TaskAssignor taskAssignor() { try { return Utils.newInstance(taskAssignorClass, TaskAssignor.class); } catch (final ClassNotFoundException e) { @@ -340,6 +319,27 @@ public TaskAssignor getTaskAssignor() { } } + public AssignmentListener assignmentListener() { + final Object o = internalConfigs.get(InternalConfig.ASSIGNMENT_LISTENER); + if (o == null) { + return stable -> { }; + } + + if (!(o instanceof AssignmentListener)) { + final KafkaException fatalException = new KafkaException( + String.format("%s is not an instance of %s", o.getClass().getName(), AssignmentListener.class.getName()) + ); + log.error(fatalException.getMessage(), fatalException); + throw fatalException; + } + + return (AssignmentListener) o; + } + + public interface AssignmentListener { + void onAssignmentComplete(final boolean stable); + } + public static class AssignmentConfigs { public final long acceptableRecoveryLag; public final int maxWarmupReplicas; diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/EosBetaUpgradeIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/EosBetaUpgradeIntegrationTest.java index b4b9dce95d18a..aac9a8ac88ce7 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/EosBetaUpgradeIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/EosBetaUpgradeIntegrationTest.java @@ -30,12 +30,15 @@ import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.common.utils.Utils; import org.apache.kafka.streams.KafkaStreams; +import org.apache.kafka.streams.KafkaStreams.State; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.StoreQueryParameters; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; +import org.apache.kafka.streams.StreamsConfig.InternalConfig; import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster; import org.apache.kafka.streams.integration.utils.IntegrationTestUtils; +import org.apache.kafka.streams.integration.utils.IntegrationTestUtils.StableAssignmentListener; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.Transformer; import org.apache.kafka.streams.kstream.TransformerSupplier; @@ -101,31 +104,6 @@ public static Collection data() { private static final int MAX_POLL_INTERVAL_MS = 100 * 1000; private static final int MAX_WAIT_TIME_MS = 60 * 1000; - private static final List> TWO_REBALANCES_STARTUP = - Collections.unmodifiableList( - Arrays.asList( - KeyValue.pair(KafkaStreams.State.CREATED, KafkaStreams.State.REBALANCING), - KeyValue.pair(KafkaStreams.State.REBALANCING, KafkaStreams.State.RUNNING), - KeyValue.pair(KafkaStreams.State.RUNNING, KafkaStreams.State.REBALANCING), - KeyValue.pair(KafkaStreams.State.REBALANCING, KafkaStreams.State.RUNNING) - ) - ); - private static final List> TWO_REBALANCES_RUNNING = - Collections.unmodifiableList( - Arrays.asList( - KeyValue.pair(KafkaStreams.State.RUNNING, KafkaStreams.State.REBALANCING), - KeyValue.pair(KafkaStreams.State.REBALANCING, KafkaStreams.State.RUNNING), - KeyValue.pair(KafkaStreams.State.RUNNING, KafkaStreams.State.REBALANCING), - KeyValue.pair(KafkaStreams.State.REBALANCING, KafkaStreams.State.RUNNING) - ) - ); - private static final List> SINGLE_REBALANCE = - Collections.unmodifiableList( - Arrays.asList( - KeyValue.pair(KafkaStreams.State.RUNNING, KafkaStreams.State.REBALANCING), - KeyValue.pair(KafkaStreams.State.REBALANCING, KafkaStreams.State.RUNNING) - ) - ); private static final List> CLOSE = Collections.unmodifiableList( Arrays.asList( @@ -160,6 +138,8 @@ public static Collection data() { private final static String MULTI_PARTITION_OUTPUT_TOPIC = "multiPartitionOutputTopic"; private final String storeName = "store"; + private final StableAssignmentListener assignmentListener = new StableAssignmentListener(); + private final AtomicBoolean errorInjectedClient1 = new AtomicBoolean(false); private final AtomicBoolean errorInjectedClient2 = new AtomicBoolean(false); private final AtomicBoolean commitErrorInjectedClient1 = new AtomicBoolean(false); @@ -261,24 +241,25 @@ public void shouldUpgradeFromEosAlphaToEosBeta() throws Exception { streams1Alpha.setStateListener( (newState, oldState) -> stateTransitions1.add(KeyValue.pair(oldState, newState)) ); + + assignmentListener.prepareForRebalance(); streams1Alpha.cleanUp(); streams1Alpha.start(); - waitForStateTransition( - stateTransitions1, - Arrays.asList( - KeyValue.pair(KafkaStreams.State.CREATED, KafkaStreams.State.REBALANCING), - KeyValue.pair(KafkaStreams.State.REBALANCING, KafkaStreams.State.RUNNING) - ) - ); + assignmentListener.waitForNextStableAssignment(MAX_WAIT_TIME_MS); + waitForRunning(stateTransitions1); - stateTransitions1.clear(); streams2Alpha = getKafkaStreams("appDir2", StreamsConfig.EXACTLY_ONCE); streams2Alpha.setStateListener( (newState, oldState) -> stateTransitions2.add(KeyValue.pair(oldState, newState)) ); + stateTransitions1.clear(); + + assignmentListener.prepareForRebalance(); streams2Alpha.cleanUp(); streams2Alpha.start(); - waitForStateTransition(stateTransitions1, TWO_REBALANCES_RUNNING, stateTransitions2, TWO_REBALANCES_STARTUP); + assignmentListener.waitForNextStableAssignment(MAX_WAIT_TIME_MS); + waitForRunning(stateTransitions1); + waitForRunning(stateTransitions2); // in all phases, we write comments that assume that p-0/p-1 are assigned to the first client // and p-2/p-3 are assigned to the second client (in reality the assignment might be different though) @@ -365,6 +346,8 @@ public void shouldUpgradeFromEosAlphaToEosBeta() throws Exception { // p-2: 10 rec + C + 5 rec (pending) // p-3: 10 rec + C + 5 rec (pending) stateTransitions2.clear(); + assignmentListener.prepareForRebalance(); + if (!injectError) { stateTransitions1.clear(); streams1Alpha.close(); @@ -377,7 +360,8 @@ public void shouldUpgradeFromEosAlphaToEosBeta() throws Exception { uncommittedInputDataBeforeFirstUpgrade.addAll(dataPotentiallyFirstFailingKey); writeInputData(dataPotentiallyFirstFailingKey); } - waitForStateTransition(stateTransitions2, SINGLE_REBALANCE); + assignmentListener.waitForNextStableAssignment(MAX_WAIT_TIME_MS); + waitForRunning(stateTransitions2); if (!injectError) { final List> committedInputDataDuringFirstUpgrade = @@ -424,8 +408,11 @@ public void shouldUpgradeFromEosAlphaToEosBeta() throws Exception { stateTransitions2.clear(); streams1Beta = getKafkaStreams("appDir1", StreamsConfig.EXACTLY_ONCE_BETA); streams1Beta.setStateListener((newState, oldState) -> stateTransitions1.add(KeyValue.pair(oldState, newState))); + assignmentListener.prepareForRebalance(); streams1Beta.start(); - waitForStateTransition(stateTransitions1, TWO_REBALANCES_STARTUP, stateTransitions2, TWO_REBALANCES_RUNNING); + assignmentListener.waitForNextStableAssignment(MAX_WAIT_TIME_MS); + waitForRunning(stateTransitions1); + waitForRunning(stateTransitions2); final Set committedKeys = mkSet(0L, 1L, 2L, 3L); if (!injectError) { @@ -497,6 +484,8 @@ public void shouldUpgradeFromEosAlphaToEosBeta() throws Exception { stateTransitions1.clear(); stateTransitions2.clear(); + assignmentListener.prepareForRebalance(); + commitCounterClient1.set(0); commitErrorInjectedClient2.set(true); @@ -509,7 +498,9 @@ public void shouldUpgradeFromEosAlphaToEosBeta() throws Exception { ); verifyUncommitted(expectedUncommittedResult); - waitForStateTransition(stateTransitions1, SINGLE_REBALANCE, stateTransitions2, CRASH); + assignmentListener.waitForNextStableAssignment(MAX_WAIT_TIME_MS); + + waitForStateTransition(stateTransitions2, CRASH); commitErrorInjectedClient2.set(false); stateTransitions2.clear(); @@ -547,8 +538,11 @@ public void shouldUpgradeFromEosAlphaToEosBeta() throws Exception { streams2AlphaTwo.setStateListener( (newState, oldState) -> stateTransitions2.add(KeyValue.pair(oldState, newState)) ); + assignmentListener.prepareForRebalance(); streams2AlphaTwo.start(); - waitForStateTransition(stateTransitions1, TWO_REBALANCES_RUNNING, stateTransitions2, TWO_REBALANCES_STARTUP); + assignmentListener.waitForNextStableAssignment(MAX_WAIT_TIME_MS); + waitForRunning(stateTransitions1); + waitForRunning(stateTransitions2); // 7b. write third batch of input data final Set keysFirstClient = keysFromInstance(streams1Beta); @@ -582,6 +576,7 @@ public void shouldUpgradeFromEosAlphaToEosBeta() throws Exception { stateTransitions1.clear(); stateTransitions2.clear(); + assignmentListener.prepareForRebalance(); commitCounterClient2.set(0); commitErrorInjectedClient1.set(true); @@ -594,7 +589,8 @@ public void shouldUpgradeFromEosAlphaToEosBeta() throws Exception { ); verifyUncommitted(expectedUncommittedResult); - waitForStateTransition(stateTransitions1, CRASH, stateTransitions2, SINGLE_REBALANCE); + assignmentListener.waitForNextStableAssignment(MAX_WAIT_TIME_MS); + waitForStateTransition(stateTransitions1, CRASH); commitErrorInjectedClient1.set(false); stateTransitions1.clear(); @@ -611,8 +607,11 @@ public void shouldUpgradeFromEosAlphaToEosBeta() throws Exception { stateTransitions2.clear(); streams1BetaTwo = getKafkaStreams("appDir1", StreamsConfig.EXACTLY_ONCE_BETA); streams1BetaTwo.setStateListener((newState, oldState) -> stateTransitions1.add(KeyValue.pair(oldState, newState))); + assignmentListener.prepareForRebalance(); streams1BetaTwo.start(); - waitForStateTransition(stateTransitions1, TWO_REBALANCES_STARTUP, stateTransitions2, TWO_REBALANCES_RUNNING); + assignmentListener.waitForNextStableAssignment(MAX_WAIT_TIME_MS); + waitForRunning(stateTransitions1); + waitForRunning(stateTransitions2); } // phase 8: (write partial fourth batch of data) @@ -675,6 +674,7 @@ public void shouldUpgradeFromEosAlphaToEosBeta() throws Exception { // p-2: 10 rec + C + 5 rec + C + 5 rec + A + 5 rec + C + 10 rec + C + 4 rec ---> A + 5 rec (pending) // p-3: 10 rec + C + 5 rec + C + 5 rec + A + 5 rec + C + 10 rec + C + 5 rec ---> A + 5 rec (pending) stateTransitions1.clear(); + assignmentListener.prepareForRebalance(); if (!injectError) { stateTransitions2.clear(); streams2AlphaTwo.close(); @@ -687,7 +687,8 @@ public void shouldUpgradeFromEosAlphaToEosBeta() throws Exception { uncommittedInputDataBeforeSecondUpgrade.addAll(dataPotentiallySecondFailingKey); writeInputData(dataPotentiallySecondFailingKey); } - waitForStateTransition(stateTransitions1, SINGLE_REBALANCE); + assignmentListener.waitForNextStableAssignment(MAX_WAIT_TIME_MS); + waitForRunning(stateTransitions1); if (!injectError) { final List> committedInputDataDuringSecondUpgrade = @@ -739,8 +740,11 @@ public void shouldUpgradeFromEosAlphaToEosBeta() throws Exception { streams2Beta.setStateListener( (newState, oldState) -> stateTransitions2.add(KeyValue.pair(oldState, newState)) ); + assignmentListener.prepareForRebalance(); streams2Beta.start(); - waitForStateTransition(stateTransitions1, TWO_REBALANCES_RUNNING, stateTransitions2, TWO_REBALANCES_STARTUP); + assignmentListener.waitForNextStableAssignment(MAX_WAIT_TIME_MS); + waitForRunning(stateTransitions1); + waitForRunning(stateTransitions2); committedKeys.addAll(mkSet(0L, 1L, 2L, 3L)); if (!injectError) { @@ -884,6 +888,7 @@ public void close() { } properties.put(StreamsConfig.producerPrefix(ProducerConfig.PARTITIONER_CLASS_CONFIG), KeyPartitioner.class); properties.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 0); properties.put(StreamsConfig.STATE_DIR_CONFIG, TestUtils.tempDirectory().getPath() + File.separator + appDir); + properties.put(InternalConfig.ASSIGNMENT_LISTENER, assignmentListener); final Properties config = StreamsTestUtils.getStreamsConfig( applicationId, @@ -906,29 +911,22 @@ public void close() { } return streams; } - private void waitForStateTransition(final List> observed, - final List> expected) - throws Exception { - + private void waitForRunning(final List> observed) throws Exception { waitForCondition( - () -> observed.equals(expected), + () -> !observed.isEmpty() && observed.get(observed.size() - 1).value.equals(State.RUNNING), MAX_WAIT_TIME_MS, () -> "Client did not startup on time. Observers transitions: " + observed ); } - private void waitForStateTransition(final List> observed1, - final List> expected1, - final List> observed2, - final List> expected2) + private void waitForStateTransition(final List> observed, + final List> expected) throws Exception { waitForCondition( - () -> observed1.equals(expected1) && observed2.equals(expected2), + () -> observed.equals(expected), MAX_WAIT_TIME_MS, - () -> "Clients did not startup and stabilize on time. Observed transitions: " + - "\n client-1 transitions: " + observed1 + - "\n client-2 transitions: " + observed2 + () -> "Client did not startup on time. Observers transitions: " + observed ); } diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/utils/IntegrationTestUtils.java b/streams/src/test/java/org/apache/kafka/streams/integration/utils/IntegrationTestUtils.java index 222d2780911f7..5c0497bccf41c 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/utils/IntegrationTestUtils.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/utils/IntegrationTestUtils.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.streams.integration.utils; +import java.util.concurrent.atomic.AtomicInteger; import kafka.api.Request; import kafka.server.KafkaServer; import kafka.server.MetadataCache; @@ -47,6 +48,7 @@ import org.apache.kafka.streams.errors.InvalidStateStoreException; import org.apache.kafka.streams.processor.internals.StreamThread; import org.apache.kafka.streams.processor.internals.ThreadStateTransitionValidator; +import org.apache.kafka.streams.processor.internals.assignment.AssignorConfiguration.AssignmentListener; import org.apache.kafka.streams.state.QueryableStoreType; import org.apache.kafka.test.TestCondition; import org.apache.kafka.test.TestUtils; @@ -82,6 +84,7 @@ import java.util.stream.Collectors; import static org.apache.kafka.test.TestUtils.retryOnExceptionWithTimeout; +import static org.apache.kafka.test.TestUtils.waitForCondition; import static org.hamcrest.MatcherAssert.assertThat; import static org.hamcrest.Matchers.equalTo; import static org.hamcrest.Matchers.greaterThanOrEqualTo; @@ -1206,4 +1209,42 @@ public static S getStore(final long waitTime, Thread.sleep(Math.min(100L, waitTime)); } } + + public static class StableAssignmentListener implements AssignmentListener { + final AtomicInteger numStableAssignments = new AtomicInteger(0); + int nextExpectedNumStableAssignments; + + @Override + public void onAssignmentComplete(final boolean stable) { + if (stable) { + numStableAssignments.incrementAndGet(); + } + } + + public int numStableAssignments() { + return numStableAssignments.get(); + } + + /** + * Saves the current number of stable rebalances so that we can tell when the next stable assignment has been + * reached. This should be called once for every invocation of {@link #waitForNextStableAssignment(long)}, + * before the rebalance-triggering event. + */ + public void prepareForRebalance() { + nextExpectedNumStableAssignments = numStableAssignments.get() + 1; + } + + /** + * Waits for the assignment to stabilize after the group rebalances. You must call {@link #prepareForRebalance()} + * prior to the rebalance-triggering event before using this method to wait. + */ + public void waitForNextStableAssignment(final long maxWaitMs) throws InterruptedException { + waitForCondition( + () -> nextExpectedNumStableAssignments == numStableAssignments(), + maxWaitMs, + () -> "Client did not reach " + nextExpectedNumStableAssignments + " stable assignments on time, " + + "numStableAssignments was " + numStableAssignments() + ); + } + } } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignorTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignorTest.java index 4c648249c54ee..644e270901424 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignorTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamsPartitionAssignorTest.java @@ -1274,6 +1274,7 @@ public void shouldThrowExceptionIfApplicationServerConfigIsNotHostPortPair() { @Test public void shouldThrowExceptionIfApplicationServerConfigPortIsNotAnInteger() { + createDefaultMockTaskManager(); assertThrows(ConfigException.class, () -> configurePartitionAssignorWith(Collections.singletonMap(StreamsConfig.APPLICATION_SERVER_CONFIG, "localhost:j87yhk"))); } @@ -1793,7 +1794,7 @@ public void shouldSetAdminClientTimeout() { props.put(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG, 2 * 60 * 1000); final AssignorConfiguration assignorConfiguration = new AssignorConfiguration(props); - assertThat(assignorConfiguration.getAdminClientTimeout(), is(2 * 60 * 1000)); + assertThat(assignorConfiguration.adminClientTimeout(), is(2 * 60 * 1000)); } @Test @@ -1804,7 +1805,7 @@ public void shouldGetNextProbingRebalanceMs() { final Map props = configProps(); final AssignorConfiguration assignorConfiguration = new AssignorConfiguration(props); - assertThat(assignorConfiguration.getNextScheduledRebalanceMs(props).get(), equalTo(5 * 60 * 1000L)); + assertThat(assignorConfiguration.nextScheduledRebalanceMs().get(), equalTo(5 * 60 * 1000L)); } @Test @@ -1815,7 +1816,7 @@ public void shouldGetTime() { final Map props = configProps(); final AssignorConfiguration assignorConfiguration = new AssignorConfiguration(props); - assertThat(assignorConfiguration.getTime(props).milliseconds(), equalTo(Long.MAX_VALUE)); + assertThat(assignorConfiguration.time().milliseconds(), equalTo(Long.MAX_VALUE)); } @Test diff --git a/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java b/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java index cb8787f7c87e6..af4d69c93ec3a 100644 --- a/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/tests/StreamsUpgradeTest.java @@ -138,7 +138,7 @@ public void configure(final Map configs) { usedSubscriptionMetadataVersionPeek = new AtomicInteger(); } configs.remove("test.future.metadata"); - nextScheduledRebalanceMs = new AssignorConfiguration(configs).getNextScheduledRebalanceMs(configs); + nextScheduledRebalanceMs = new AssignorConfiguration(configs).nextScheduledRebalanceMs(); super.configure(configs); }