You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@kafka.apache.org by gu...@apache.org on 2022/03/23 03:56:07 UTC
[kafka] branch trunk updated: [Emit final][4/N] add time ordered store factory (#11892)
This is an automated email from the ASF dual-hosted git repository.
guozhang pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new a3adf41 [Emit final][4/N] add time ordered store factory (#11892)
a3adf41 is described below
commit a3adf41d8b90d1244232e448b959db3b3f4dc2fe
Author: Hao Li <11...@users.noreply.github.com>
AuthorDate: Tue Mar 22 20:53:53 2022 -0700
[Emit final][4/N] add time ordered store factory (#11892)
Add factory to create time ordered store supplier.
Reviewers: Guozhang Wang <wa...@gmail.com>
---
...IndexedTimeOrderedWindowBytesStoreSupplier.java | 37 +++++++++++
...xedTimeOrderedWindowBytesStoreSupplierTest.java | 75 ++++++++++++++++++++++
2 files changed, 112 insertions(+)
diff --git a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDbIndexedTimeOrderedWindowBytesStoreSupplier.java b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDbIndexedTimeOrderedWindowBytesStoreSupplier.java
index 84d8a80..af5417f 100644
--- a/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDbIndexedTimeOrderedWindowBytesStoreSupplier.java
+++ b/streams/src/main/java/org/apache/kafka/streams/state/internals/RocksDbIndexedTimeOrderedWindowBytesStoreSupplier.java
@@ -16,6 +16,11 @@
*/
package org.apache.kafka.streams.state.internals;
+import static org.apache.kafka.streams.internals.ApiUtils.prepareMillisCheckFailMsgPrefix;
+import static org.apache.kafka.streams.internals.ApiUtils.validateMillisecondDuration;
+
+import java.time.Duration;
+import java.util.Objects;
import org.apache.kafka.common.utils.Bytes;
import org.apache.kafka.streams.state.WindowBytesStoreSupplier;
import org.apache.kafka.streams.state.WindowStore;
@@ -33,6 +38,38 @@ public class RocksDbIndexedTimeOrderedWindowBytesStoreSupplier implements Window
private final boolean retainDuplicates;
private final WindowStoreTypes windowStoreType;
+ public static RocksDbIndexedTimeOrderedWindowBytesStoreSupplier create(final String name,
+ final Duration retentionPeriod,
+ final Duration windowSize,
+ final boolean retainDuplicates,
+ final boolean hasIndex) {
+ Objects.requireNonNull(name, "name cannot be null");
+ final String rpMsgPrefix = prepareMillisCheckFailMsgPrefix(retentionPeriod, "retentionPeriod");
+ final long retentionMs = validateMillisecondDuration(retentionPeriod, rpMsgPrefix);
+ final String wsMsgPrefix = prepareMillisCheckFailMsgPrefix(windowSize, "windowSize");
+ final long windowSizeMs = validateMillisecondDuration(windowSize, wsMsgPrefix);
+
+ final long defaultSegmentInterval = Math.max(retentionMs / 2, 60_000L);
+
+ if (retentionMs < 0L) {
+ throw new IllegalArgumentException("retentionPeriod cannot be negative");
+ }
+ if (windowSizeMs < 0L) {
+ throw new IllegalArgumentException("windowSize cannot be negative");
+ }
+ if (defaultSegmentInterval < 1L) {
+ throw new IllegalArgumentException("segmentInterval cannot be zero or negative");
+ }
+ if (windowSizeMs > retentionMs) {
+ throw new IllegalArgumentException("The retention period of the window store "
+ + name + " must be no smaller than its window size. Got size=["
+ + windowSizeMs + "], retention=[" + retentionMs + "]");
+ }
+
+ return new RocksDbIndexedTimeOrderedWindowBytesStoreSupplier(name, retentionMs,
+ defaultSegmentInterval, windowSizeMs, retainDuplicates, hasIndex);
+ }
+
public RocksDbIndexedTimeOrderedWindowBytesStoreSupplier(final String name,
final long retentionPeriod,
final long segmentInterval,
diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/RocksDbIndexedTimeOrderedWindowBytesStoreSupplierTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/RocksDbIndexedTimeOrderedWindowBytesStoreSupplierTest.java
new file mode 100644
index 0000000..ed1bbb8
--- /dev/null
+++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/RocksDbIndexedTimeOrderedWindowBytesStoreSupplierTest.java
@@ -0,0 +1,75 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.kafka.streams.state.internals;
+
+import org.apache.kafka.streams.processor.StateStore;
+import org.apache.kafka.streams.state.WindowStore;
+import org.junit.Test;
+
+import static java.time.Duration.ZERO;
+import static java.time.Duration.ofMillis;
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.core.IsInstanceOf.instanceOf;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertThrows;
+import static org.junit.Assert.assertTrue;
+
+public class RocksDbIndexedTimeOrderedWindowBytesStoreSupplierTest {
+
+ @Test
+ public void shouldThrowIfStoreNameIsNull() {
+ final Exception e = assertThrows(NullPointerException.class, () -> RocksDbIndexedTimeOrderedWindowBytesStoreSupplier.create(null, ZERO, ZERO, false, false));
+ assertEquals("name cannot be null", e.getMessage());
+ }
+
+ @Test
+ public void shouldThrowIfRetentionPeriodIsNegative() {
+ final Exception e = assertThrows(IllegalArgumentException.class, () -> RocksDbIndexedTimeOrderedWindowBytesStoreSupplier.create("anyName", ofMillis(-1L), ZERO, false, false));
+ assertEquals("retentionPeriod cannot be negative", e.getMessage());
+ }
+
+ @Test
+ public void shouldThrowIfWindowSizeIsNegative() {
+ final Exception e = assertThrows(IllegalArgumentException.class, () -> RocksDbIndexedTimeOrderedWindowBytesStoreSupplier.create("anyName", ofMillis(0L), ofMillis(-1L), false, false));
+ assertEquals("windowSize cannot be negative", e.getMessage());
+ }
+
+ @Test
+ public void shouldThrowIfWindowSizeIsLargerThanRetention() {
+ final Exception e = assertThrows(IllegalArgumentException.class, () -> RocksDbIndexedTimeOrderedWindowBytesStoreSupplier.create("anyName", ofMillis(1L), ofMillis(2L), false, false));
+ assertEquals("The retention period of the window store anyName must be no smaller than its window size. Got size=[2], retention=[1]", e.getMessage());
+ }
+
+ @Test
+ public void shouldCreateRocksDbTimeOrderedWindowStoreWithIndex() {
+ final WindowStore store = RocksDbIndexedTimeOrderedWindowBytesStoreSupplier.create("store", ofMillis(1L), ofMillis(1L), false, true).get();
+ final StateStore wrapped = ((WrappedStateStore) store).wrapped();
+ assertThat(store, instanceOf(RocksDBTimeOrderedWindowStore.class));
+ assertThat(wrapped, instanceOf(RocksDBTimeOrderedSegmentedBytesStore.class));
+ assertTrue(((RocksDBTimeOrderedSegmentedBytesStore) wrapped).hasIndex());
+ }
+
+ @Test
+ public void shouldCreateRocksDbTimeOrderedWindowStoreWithoutIndex() {
+ final WindowStore store = RocksDbIndexedTimeOrderedWindowBytesStoreSupplier.create("store", ofMillis(1L), ofMillis(1L), false, false).get();
+ final StateStore wrapped = ((WrappedStateStore) store).wrapped();
+ assertThat(store, instanceOf(RocksDBTimeOrderedWindowStore.class));
+ assertThat(wrapped, instanceOf(RocksDBTimeOrderedSegmentedBytesStore.class));
+ assertFalse(((RocksDBTimeOrderedSegmentedBytesStore) wrapped).hasIndex());
+ }
+}