From d16a0f5c790a0b1815688873775dac6ca3601376 Mon Sep 17 00:00:00 2001 From: Tomoyuki Saito Date: Wed, 17 Oct 2018 10:08:09 +0900 Subject: [PATCH 1/4] MINOR: Prohibit setting StreamsConfig commit.interval.ms to a negative value --- .../src/main/java/org/apache/kafka/streams/StreamsConfig.java | 1 + .../apache/kafka/streams/integration/EosIntegrationTest.java | 2 +- 2 files changed, 2 insertions(+), 1 deletion(-) 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 b112a54d1f791..b25894c2e0452 100644 --- a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java +++ b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java @@ -582,6 +582,7 @@ public class StreamsConfig extends AbstractConfig { .define(COMMIT_INTERVAL_MS_CONFIG, Type.LONG, DEFAULT_COMMIT_INTERVAL_MS, + atLeast(0), Importance.LOW, COMMIT_INTERVAL_MS_DOC) .define(CONNECTIONS_MAX_IDLE_MS_CONFIG, diff --git a/streams/src/test/java/org/apache/kafka/streams/integration/EosIntegrationTest.java b/streams/src/test/java/org/apache/kafka/streams/integration/EosIntegrationTest.java index 9aa07729c2075..b2bd0a87771e1 100644 --- a/streams/src/test/java/org/apache/kafka/streams/integration/EosIntegrationTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/integration/EosIntegrationTest.java @@ -677,7 +677,7 @@ public void close() { } { put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE); put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, numberOfStreamsThreads); - put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, -1); + put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, Long.MAX_VALUE); put(StreamsConfig.consumerPrefix(ConsumerConfig.METADATA_MAX_AGE_CONFIG), "1000"); put(StreamsConfig.consumerPrefix(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG), "earliest"); put(StreamsConfig.consumerPrefix(ConsumerConfig.REQUEST_TIMEOUT_MS_CONFIG), 5 * 1000); From 627a938cf96d998f4920faf53c86967cab451745 Mon Sep 17 00:00:00 2001 From: Tomoyuki Saito Date: Fri, 19 Oct 2018 07:46:19 +0900 Subject: [PATCH 2/4] (fixup) Remove the check interval >= 0 --- .../kafka/streams/processor/internals/GlobalStreamThread.java | 2 +- .../apache/kafka/streams/processor/internals/StreamThread.java | 2 +- 2 files changed, 2 insertions(+), 2 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 9d529c5455c46..d91aedf6c3eed 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 @@ -241,7 +241,7 @@ void pollAndUpdate() { stateMaintainer.update(record); } final long now = time.milliseconds(); - if (flushInterval >= 0 && now >= lastFlush + flushInterval) { + if (now >= lastFlush + flushInterval) { stateMaintainer.flushState(); lastFlush = now; } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index 3a19aa7910cbb..6c3e7cd8e7b91 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -1020,7 +1020,7 @@ private boolean maybePunctuate() { boolean maybeCommit() { int committed = 0; - if (commitTimeMs >= 0 && now - lastCommitMs > commitTimeMs) { + if (now - lastCommitMs > commitTimeMs) { if (log.isTraceEnabled()) { log.trace("Committing all active tasks {} and standby tasks {} since {}ms has elapsed (commit interval is {}ms)", taskManager.activeTaskIds(), taskManager.standbyTaskIds(), now - lastCommitMs, commitTimeMs); From eb7f2cdafba055ccaaee9300a53a4a360e1ab641 Mon Sep 17 00:00:00 2001 From: Tomoyuki Saito Date: Fri, 19 Oct 2018 09:54:02 +0900 Subject: [PATCH 3/4] (fixup) Add a test --- .../org/apache/kafka/streams/StreamsConfigTest.java | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java b/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java index 83279cc4ddca5..a86c38946df77 100644 --- a/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/StreamsConfigTest.java @@ -521,6 +521,18 @@ public void shouldNotOverrideUserConfigCommitIntervalMsIfExactlyOnceEnabled() { assertThat(streamsConfig.getLong(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG), equalTo(commitIntervalMs)); } + @Test + public void shouldThrowExceptionIfCommitIntervalMsIsNegative() { + final long commitIntervalMs = -1; + props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, commitIntervalMs); + try { + new StreamsConfig(props); + fail("Should throw ConfigException when commitIntervalMs is set to a negative value"); + } catch (final ConfigException e) { + assertEquals("Invalid value -1 for configuration commit.interval.ms: Value must be at least 0", e.getMessage()); + } + } + @Test public void shouldUseNewConfigsWhenPresent() { final Properties props = getStreamsConfig(); From 35e51dcfd8e9c881ff8e5c950ed8be42b4aa5f6a Mon Sep 17 00:00:00 2001 From: Tomoyuki Saito Date: Fri, 19 Oct 2018 10:13:18 +0900 Subject: [PATCH 4/4] (fixup) Remove StateConsumer spec when setting flushInterval to a negative value --- .../streams/processor/internals/StateConsumerTest.java | 9 --------- 1 file changed, 9 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateConsumerTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateConsumerTest.java index 140f705619956..f44d2b45ecf16 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateConsumerTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StateConsumerTest.java @@ -108,15 +108,6 @@ public void shouldNotFlushOffsetsWhenFlushIntervalHasNotLapsed() { assertFalse(stateMaintainer.flushed); } - @Test - public void shouldNotFlushWhenFlushIntervalIsZero() { - stateConsumer = new GlobalStreamThread.StateConsumer(logContext, consumer, stateMaintainer, time, Duration.ofMillis(10L), -1); - stateConsumer.initialize(); - time.sleep(100); - stateConsumer.pollAndUpdate(); - assertFalse(stateMaintainer.flushed); - } - @Test public void shouldCloseConsumer() throws IOException { stateConsumer.close();