Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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 @@ -48,7 +48,6 @@
import java.util.Collections;
import java.util.List;
import java.util.Properties;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.atomic.AtomicBoolean;

import static org.apache.kafka.common.utils.Utils.mkEntry;
Expand All @@ -68,6 +67,7 @@
public class StreamsUncaughtExceptionHandlerIntegrationTest {
@ClassRule
public static final EmbeddedKafkaCluster CLUSTER = new EmbeddedKafkaCluster(1);
public static final Duration DEFAULT_DURATION = Duration.ofSeconds(30);

@Rule
public TestName testName = new TestName();
Expand Down Expand Up @@ -110,83 +110,45 @@ public void teardown() throws IOException {
}

@Test
public void shouldShutdownThreadUsingOldHandler() throws Exception {
public void shouldShutdownThreadUsingOldHandler() throws InterruptedException {
try (final KafkaStreams kafkaStreams = new KafkaStreams(builder.build(), properties)) {
final CountDownLatch latch = new CountDownLatch(1);

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

delete unused latch

final AtomicBoolean flag = new AtomicBoolean(false);
kafkaStreams.setUncaughtExceptionHandler((t, e) -> flag.set(true));

StreamsTestUtils.startKafkaStreamsAndWaitForRunningState(kafkaStreams);

produceMessages(0L, inputTopic, "A");
waitForApplicationState(Collections.singletonList(kafkaStreams), KafkaStreams.State.ERROR, Duration.ofSeconds(15));

TestUtils.waitForCondition(flag::get, "Handler was called");
waitForApplicationState(Collections.singletonList(kafkaStreams), KafkaStreams.State.ERROR, DEFAULT_DURATION);

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We should wait for the uncaughtExceptionHandler got called before waiting for the streams state change.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The order is not really that important here, either way works


assertThat(processorValueCollector.size(), equalTo(2));

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@wcarlson5 for example, this test probably should have multiple threads, right?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

as above

}
}

@Test
public void shouldShutdownClient() throws Exception {
public void shouldShutdownClient() throws InterruptedException {
try (final KafkaStreams kafkaStreams = new KafkaStreams(builder.build(), properties)) {
final CountDownLatch latch = new CountDownLatch(1);
kafkaStreams.setUncaughtExceptionHandler((t, e) -> fail("should not hit old handler"));

kafkaStreams.setUncaughtExceptionHandler(exception -> SHUTDOWN_CLIENT);

StreamsTestUtils.startKafkaStreamsAndWaitForRunningState(kafkaStreams);

produceMessages(0L, inputTopic, "A");
waitForApplicationState(Collections.singletonList(kafkaStreams), KafkaStreams.State.NOT_RUNNING, Duration.ofSeconds(15));
waitForApplicationState(Collections.singletonList(kafkaStreams), KafkaStreams.State.NOT_RUNNING, DEFAULT_DURATION);

assertThat(processorValueCollector.size(), equalTo(1));
}
}

@Test
public void shouldShutdownApplication() throws Exception {
final Topology topology = builder.build();

try (final KafkaStreams kafkaStreams = new KafkaStreams(topology, properties)) {
final KafkaStreams kafkaStreams1 = new KafkaStreams(topology, properties);

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

forgot to put the kafkaStreams1 into the try resource block here

final CountDownLatch latch = new CountDownLatch(1);
kafkaStreams.setUncaughtExceptionHandler((t, e) -> fail("should not hit old handler"));
kafkaStreams1.setUncaughtExceptionHandler((t, e) -> fail("should not hit old handler"));
kafkaStreams.setUncaughtExceptionHandler(exception -> SHUTDOWN_APPLICATION);
kafkaStreams1.setUncaughtExceptionHandler(exception -> SHUTDOWN_APPLICATION);

kafkaStreams.start();
kafkaStreams1.start();

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The start is async, and we didn't wait for it.


produceMessages(0L, inputTopic, "A");
waitForApplicationState(Arrays.asList(kafkaStreams, kafkaStreams1), KafkaStreams.State.ERROR, Duration.ofSeconds(30));

assertThat(processorValueCollector.size(), equalTo(1));
}
testShutdownApplication(2);
}

@Test
public void shouldShutdownSingleThreadApplication() throws Exception {
properties.setProperty(StreamsConfig.NUM_STREAM_THREADS_CONFIG, "1");

final Topology topology = builder.build();

try (final KafkaStreams kafkaStreams = new KafkaStreams(topology, properties)) {
final KafkaStreams kafkaStreams1 = new KafkaStreams(topology, properties);
final CountDownLatch latch = new CountDownLatch(1);
kafkaStreams.setUncaughtExceptionHandler((t, e) -> fail("should not hit old handler"));
kafkaStreams1.setUncaughtExceptionHandler((t, e) -> fail("should not hit old handler"));
kafkaStreams.setUncaughtExceptionHandler(exception -> SHUTDOWN_APPLICATION);
kafkaStreams1.setUncaughtExceptionHandler(exception -> SHUTDOWN_APPLICATION);

kafkaStreams.start();
kafkaStreams1.start();

produceMessages(0L, inputTopic, "A");
waitForApplicationState(Arrays.asList(kafkaStreams, kafkaStreams1), KafkaStreams.State.ERROR, Duration.ofSeconds(30));

assertThat(processorValueCollector.size(), equalTo(1));
}
testShutdownApplication(1);
}

private void produceMessages(final long timestamp, final String streamOneInput, final String msg) {
Expand All @@ -202,7 +164,6 @@ private void produceMessages(final long timestamp, final String streamOneInput,
}

private static class ShutdownProcessor extends AbstractProcessor<String, String> {

final List<String> valueList;

ShutdownProcessor(final List<String> valueList) {
Expand All @@ -215,6 +176,28 @@ public void process(final String key, final String value) {
throw new StreamsException(Thread.currentThread().getName());
}
}

private void testShutdownApplication(final int numThreads) throws InterruptedException {
properties.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, numThreads);

final Topology topology = builder.build();

try (final KafkaStreams kafkaStreams1 = new KafkaStreams(topology, properties);
final KafkaStreams kafkaStreams2 = new KafkaStreams(topology, properties)) {
kafkaStreams1.setUncaughtExceptionHandler((t, e) -> fail("should not hit old handler"));
kafkaStreams2.setUncaughtExceptionHandler((t, e) -> fail("should not hit old handler"));
kafkaStreams1.setUncaughtExceptionHandler(exception -> SHUTDOWN_APPLICATION);
kafkaStreams2.setUncaughtExceptionHandler(exception -> SHUTDOWN_APPLICATION);

StreamsTestUtils.startKafkaStreamsAndWaitForRunningState(kafkaStreams1);
StreamsTestUtils.startKafkaStreamsAndWaitForRunningState(kafkaStreams2);

produceMessages(0L, inputTopic, "A");
waitForApplicationState(Arrays.asList(kafkaStreams1, kafkaStreams2), KafkaStreams.State.ERROR, DEFAULT_DURATION);

assertThat(processorValueCollector.size(), equalTo(1));
}
}
}


Original file line number Diff line number Diff line change
Expand Up @@ -833,7 +833,7 @@ private static void waitUntilMetadataIsPropagated(final List<KafkaServer> server
* {@link State#RUNNING} state at the same time. Note that states may change between the time
* that this method returns and the calling function executes its next statement.<p>
*
* When the application is already started use {@link #waitForApplicationState(List, State, Duration)}
* If the application is already started, use {@link #waitForApplicationState(List, State, Duration)}
* to wait for instances to reach {@link State#RUNNING} state.
*
* @param streamsList the list of streams instances to run.
Expand Down Expand Up @@ -903,16 +903,19 @@ public static void startApplicationAndWaitUntilRunning(final List<KafkaStreams>
}

/**
* Waits for the given {@link KafkaStreams} instances to all be in a {@link State#RUNNING}
* state. Prefer {@link #startApplicationAndWaitUntilRunning(List, Duration)} when possible
* Waits for the given {@link KafkaStreams} instances to all be in a specific {@link State}.
* Prefer {@link #startApplicationAndWaitUntilRunning(List, Duration)} when possible
* because this method uses polling, which can be more error prone and slightly slower.
*
* @param streamsList the list of streams instances to run.
* @param timeout the time to wait for the streams to all be in {@link State#RUNNING} state.
* @param state the expected state that all the streams to be in within timeout
* @param timeout the time to wait for the streams to all be in the specific state.
*
* @throws InterruptedException if the streams doesn't change to the expected state in time.
*/
public static void waitForApplicationState(final List<KafkaStreams> streamsList,
final State state,
final Duration timeout) throws Exception {
final Duration timeout) throws InterruptedException {

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We should throw a specific kind of exception, not an Exception.

retryOnExceptionWithTimeout(timeout.toMillis(), () -> {
final Map<KafkaStreams, State> streamsToStates = streamsList
.stream()
Expand Down