From e56a1a6f593ac85df8f11191f9d6e6ec0b0b35c2 Mon Sep 17 00:00:00 2001 From: Cameron Lee Date: Fri, 24 Apr 2020 14:58:25 -0700 Subject: [PATCH 1/2] SAMZA-2516: Migrate BaseKeyValueStorageEngineFactory to be an abstract class instead of trait --- build.gradle | 1 + .../kv/BaseKeyValueStorageEngineFactory.java | 221 ++++++++++++++ .../storage/kv/LargeMessageSafeStore.java | 6 + .../storage/kv/RecordTooLargeException.java | 0 .../samza/storage/kv/AccessLoggedStore.scala | 6 + .../kv/BaseKeyValueStorageEngineFactory.scala | 200 ------------- .../apache/samza/storage/kv/CachedStore.scala | 6 + .../storage/kv/KeyValueStorageEngine.scala | 11 + .../apache/samza/storage/kv/LoggedStore.scala | 6 + .../storage/kv/NullSafeKeyValueStore.scala | 6 + .../storage/kv/SerializedKeyValueStore.scala | 6 + .../kv/MockKeyValueStorageEngineFactory.java | 45 +++ .../TestBaseKeyValueStorageEngineFactory.java | 273 ++++++++++++++++++ 13 files changed, 587 insertions(+), 200 deletions(-) create mode 100644 samza-kv/src/main/java/org/apache/samza/storage/kv/BaseKeyValueStorageEngineFactory.java rename samza-kv/src/main/{scala => java}/org/apache/samza/storage/kv/RecordTooLargeException.java (100%) delete mode 100644 samza-kv/src/main/scala/org/apache/samza/storage/kv/BaseKeyValueStorageEngineFactory.scala create mode 100644 samza-kv/src/test/java/org/apache/samza/storage/kv/MockKeyValueStorageEngineFactory.java create mode 100644 samza-kv/src/test/java/org/apache/samza/storage/kv/TestBaseKeyValueStorageEngineFactory.java diff --git a/build.gradle b/build.gradle index 94d353b9f9..dc3c623203 100644 --- a/build.gradle +++ b/build.gradle @@ -624,6 +624,7 @@ project(":samza-kv_$scalaSuffix") { compile project(':samza-api') compile project(":samza-core_$scalaSuffix") compile "org.scala-lang:scala-library:$scalaVersion" + testCompile "com.google.guava:guava:$guavaVersion" testCompile "junit:junit:$junitVersion" testCompile "org.mockito:mockito-core:$mockitoVersion" } diff --git a/samza-kv/src/main/java/org/apache/samza/storage/kv/BaseKeyValueStorageEngineFactory.java b/samza-kv/src/main/java/org/apache/samza/storage/kv/BaseKeyValueStorageEngineFactory.java new file mode 100644 index 0000000000..44ada3cede --- /dev/null +++ b/samza-kv/src/main/java/org/apache/samza/storage/kv/BaseKeyValueStorageEngineFactory.java @@ -0,0 +1,221 @@ +/* + * 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.samza.storage.kv; + +import java.io.File; +import java.util.Optional; +import org.apache.commons.lang3.StringUtils; +import org.apache.samza.SamzaException; +import org.apache.samza.config.Config; +import org.apache.samza.config.MetricsConfig; +import org.apache.samza.config.StorageConfig; +import org.apache.samza.context.ContainerContext; +import org.apache.samza.context.JobContext; +import org.apache.samza.metrics.MetricsRegistry; +import org.apache.samza.serializers.Serde; +import org.apache.samza.storage.StorageEngine; +import org.apache.samza.storage.StorageEngineFactory; +import org.apache.samza.storage.StoreProperties; +import org.apache.samza.system.SystemStreamPartition; +import org.apache.samza.task.MessageCollector; +import org.apache.samza.util.HighResolutionClock; +import org.apache.samza.util.ScalaJavaUtil; + + +/** + * This encapsulates all the steps needed to create a key value storage engine. + * This is meant to be extended by the specific key value store factory implementations which will in turn override the + * getKVStore method to return a raw key-value store. + */ +public abstract class BaseKeyValueStorageEngineFactory implements StorageEngineFactory { + private static final String INMEMORY_KV_STORAGE_ENGINE_FACTORY = + "org.apache.samza.storage.kv.inmemory.InMemoryKeyValueStorageEngineFactory"; + + /** + * Implement this to return a KeyValueStore instance for the given store name, which will be used as the underlying + * raw store. + * + * @param storeName Name of the store + * @param storeDir The directory of the store + * @param registry MetricsRegistry to which to publish store specific metrics. + * @param changeLogSystemStreamPartition Samza stream partition from which to receive the changelog. + * @param jobContext Information about the job in which the task is executing. + * @param containerContext Information about the container in which the task is executing. + * @return A raw KeyValueStore instance + */ + protected abstract KeyValueStore getKVStore(String storeName, + File storeDir, + MetricsRegistry registry, + SystemStreamPartition changeLogSystemStreamPartition, + JobContext jobContext, + ContainerContext containerContext, + StoreMode storeMode); + + /** + * Constructs a key-value StorageEngine and returns it to the caller + * + * @param storeName The name of the storage engine. + * @param storeDir The directory of the storage engine. + * @param keySerde The serializer to use for serializing keys when reading or writing to the store. + * @param msgSerde The serializer to use for serializing messages when reading or writing to the store. + * @param changelogCollector MessageCollector the storage engine uses to persist changes. + * @param registry MetricsRegistry to which to publish storage-engine specific metrics. + * @param changelogSSP Samza system stream partition from which to receive the changelog. + * @param containerContext Information about the container in which the task is executing. + **/ + public StorageEngine getStorageEngine(String storeName, + File storeDir, + Serde keySerde, + Serde msgSerde, + MessageCollector changelogCollector, + MetricsRegistry registry, + SystemStreamPartition changelogSSP, + JobContext jobContext, + ContainerContext containerContext, + StoreMode storeMode) { + Config storageConfigSubset = jobContext.getConfig().subset("stores." + storeName + ".", true); + StorageConfig storageConfig = new StorageConfig(jobContext.getConfig()); + Optional storeFactory = storageConfig.getStorageFactoryClassName(storeName); + StoreProperties.StorePropertiesBuilder storePropertiesBuilder = new StoreProperties.StorePropertiesBuilder(); + if (!storeFactory.isPresent() || StringUtils.isBlank(storeFactory.get())) { + throw new SamzaException("Store factory not defined. Cannot proceed with KV store creation!"); + } + if (!storeFactory.get().equals(INMEMORY_KV_STORAGE_ENGINE_FACTORY)) { + storePropertiesBuilder.setPersistedToDisk(true); + } + int batchSize = storageConfigSubset.getInt("write.batch.size", 500); + int cacheSize = storageConfigSubset.getInt("object.cache.size", Math.max(batchSize, 1000)); + if (cacheSize > 0 && cacheSize < batchSize) { + throw new SamzaException( + "A store's cache.size cannot be less than batch.size as batched values reside in cache."); + } + if (keySerde == null) { + throw new SamzaException("Must define a key serde when using key value storage."); + } + if (msgSerde == null) { + throw new SamzaException("Must define a message serde when using key value storage."); + } + + KeyValueStore rawStore = + getKVStore(storeName, storeDir, registry, changelogSSP, jobContext, containerContext, storeMode); + KeyValueStore maybeLoggedStore = buildMaybeLoggedStore(changelogSSP, + storeName, registry, storePropertiesBuilder, rawStore, changelogCollector); + // this also applies serialization and caching layers + KeyValueStore toBeAccessLoggedStore = applyLargeMessageHandling(storeName, registry, + maybeLoggedStore, storageConfig, cacheSize, batchSize, keySerde, msgSerde); + KeyValueStore maybeAccessLoggedStore = + buildMaybeAccessLoggedStore(storeName, toBeAccessLoggedStore, changelogCollector, changelogSSP, storageConfig, + keySerde); + KeyValueStore nullSafeStore = new NullSafeKeyValueStore<>(maybeAccessLoggedStore); + + KeyValueStorageEngineMetrics keyValueStorageEngineMetrics = new KeyValueStorageEngineMetrics(storeName, registry); + HighResolutionClock clock = buildClock(jobContext.getConfig()); + return new KeyValueStorageEngine<>(storeName, storeDir, storePropertiesBuilder.build(), nullSafeStore, rawStore, + changelogSSP, changelogCollector, keyValueStorageEngineMetrics, batchSize, + ScalaJavaUtil.toScalaFunction(clock::nanoTime)); + } + + private static KeyValueStore buildMaybeLoggedStore(SystemStreamPartition changelogSSP, + String storeName, + MetricsRegistry registry, + StoreProperties.StorePropertiesBuilder storePropertiesBuilder, + KeyValueStore storeToWrap, + MessageCollector changelogCollector) { + if (changelogSSP == null) { + return storeToWrap; + } else { + LoggedStoreMetrics loggedStoreMetrics = new LoggedStoreMetrics(storeName, registry); + storePropertiesBuilder.setLoggedStore(true); + return new LoggedStore<>(storeToWrap, changelogSSP, changelogCollector, loggedStoreMetrics); + } + } + + private static KeyValueStore applyLargeMessageHandling(String storeName, + MetricsRegistry registry, + KeyValueStore storeToWrap, + StorageConfig storageConfig, + int cacheSize, + int batchSize, + Serde keySerde, + Serde msgSerde) { + int maxMessageSize = storageConfig.getChangelogMaxMsgSizeBytes(storeName); + if (storageConfig.getDisallowLargeMessages(storeName)) { + /* + * If large messages are disallowed in config, then this creates a LargeMessageSafeKeyValueStore that throws a + * RecordTooLargeException when a large message is encountered. + */ + KeyValueStore maybeCachedStore = + buildMaybeCachedStore(storeName, registry, storeToWrap, cacheSize, batchSize); + LargeMessageSafeStore largeMessageSafeKeyValueStore = + new LargeMessageSafeStore(maybeCachedStore, storeName, false, maxMessageSize); + return buildSerializedStore(storeName, registry, largeMessageSafeKeyValueStore, keySerde, msgSerde); + } else { + KeyValueStore toBeSerializedStore; + if (storageConfig.getDropLargeMessages(storeName)) { + toBeSerializedStore = new LargeMessageSafeStore(storeToWrap, storeName, true, maxMessageSize); + } else { + toBeSerializedStore = storeToWrap; + } + KeyValueStore serializedStore = + buildSerializedStore(storeName, registry, toBeSerializedStore, keySerde, msgSerde); + return buildMaybeCachedStore(storeName, registry, serializedStore, cacheSize, batchSize); + } + } + + private static KeyValueStore buildMaybeCachedStore(String storeName, MetricsRegistry registry, + KeyValueStore storeToWrap, int cacheSize, int batchSize) { + if (cacheSize > 0) { + CachedStoreMetrics cachedStoreMetrics = new CachedStoreMetrics(storeName, registry); + return new CachedStore<>(storeToWrap, cacheSize, batchSize, cachedStoreMetrics); + } else { + return storeToWrap; + } + } + + private static KeyValueStore buildSerializedStore(String storeName, + MetricsRegistry registry, + KeyValueStore storeToWrap, + Serde keySerde, + Serde msgSerde) { + SerializedKeyValueStoreMetrics serializedMetrics = new SerializedKeyValueStoreMetrics(storeName, registry); + return new SerializedKeyValueStore<>(storeToWrap, keySerde, msgSerde, serializedMetrics); + } + + private static KeyValueStore buildMaybeAccessLoggedStore(String storeName, + KeyValueStore storeToWrap, + MessageCollector changelogCollector, + SystemStreamPartition changelogSSP, + StorageConfig storageConfig, + Serde keySerde) { + if (storageConfig.getAccessLogEnabled(storeName)) { + return new AccessLoggedStore<>(storeToWrap, changelogCollector, changelogSSP, storageConfig, storeName, keySerde); + } else { + return storeToWrap; + } + } + + private static HighResolutionClock buildClock(Config config) { + MetricsConfig metricsConfig = new MetricsConfig(config); + if (metricsConfig.getMetricsTimerEnabled()) { + return System::nanoTime; + } else { + return () -> 0; + } + } +} diff --git a/samza-kv/src/main/java/org/apache/samza/storage/kv/LargeMessageSafeStore.java b/samza-kv/src/main/java/org/apache/samza/storage/kv/LargeMessageSafeStore.java index 177a9860bd..683cb52bd7 100644 --- a/samza-kv/src/main/java/org/apache/samza/storage/kv/LargeMessageSafeStore.java +++ b/samza-kv/src/main/java/org/apache/samza/storage/kv/LargeMessageSafeStore.java @@ -22,6 +22,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Optional; +import com.google.common.annotations.VisibleForTesting; import org.apache.samza.checkpoint.CheckpointId; import org.apache.samza.metrics.MetricsRegistryMap; import org.slf4j.Logger; @@ -155,4 +156,9 @@ private List> removeLargeMessages(List getStore() { + return this.store; + } } diff --git a/samza-kv/src/main/scala/org/apache/samza/storage/kv/RecordTooLargeException.java b/samza-kv/src/main/java/org/apache/samza/storage/kv/RecordTooLargeException.java similarity index 100% rename from samza-kv/src/main/scala/org/apache/samza/storage/kv/RecordTooLargeException.java rename to samza-kv/src/main/java/org/apache/samza/storage/kv/RecordTooLargeException.java diff --git a/samza-kv/src/main/scala/org/apache/samza/storage/kv/AccessLoggedStore.scala b/samza-kv/src/main/scala/org/apache/samza/storage/kv/AccessLoggedStore.scala index 8c32793550..a9db0827bb 100644 --- a/samza-kv/src/main/scala/org/apache/samza/storage/kv/AccessLoggedStore.scala +++ b/samza-kv/src/main/scala/org/apache/samza/storage/kv/AccessLoggedStore.scala @@ -24,6 +24,7 @@ import java.nio.file.Path import java.util import java.util.Optional +import com.google.common.annotations.VisibleForTesting import org.apache.samza.checkpoint.CheckpointId import org.apache.samza.config.StorageConfig import org.apache.samza.task.MessageCollector @@ -167,4 +168,9 @@ class AccessLoggedStore[K, V]( override def checkpoint(id: CheckpointId): Optional[Path] = { store.checkpoint(id) } + + @VisibleForTesting + private[kv] def getStore: KeyValueStore[K, V] = { + store + } } diff --git a/samza-kv/src/main/scala/org/apache/samza/storage/kv/BaseKeyValueStorageEngineFactory.scala b/samza-kv/src/main/scala/org/apache/samza/storage/kv/BaseKeyValueStorageEngineFactory.scala deleted file mode 100644 index fef1debe33..0000000000 --- a/samza-kv/src/main/scala/org/apache/samza/storage/kv/BaseKeyValueStorageEngineFactory.scala +++ /dev/null @@ -1,200 +0,0 @@ -/* - * 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.samza.storage.kv - -import java.io.File - -import org.apache.samza.SamzaException -import org.apache.samza.config.{MetricsConfig, StorageConfig} -import org.apache.samza.context.{ContainerContext, JobContext} -import org.apache.samza.metrics.MetricsRegistry -import org.apache.samza.serializers.Serde -import org.apache.samza.storage.StorageEngineFactory.StoreMode -import org.apache.samza.storage.{StorageEngine, StorageEngineFactory, StoreProperties} -import org.apache.samza.system.SystemStreamPartition -import org.apache.samza.task.MessageCollector -import org.apache.samza.util.ScalaJavaUtil.JavaOptionals -import org.apache.samza.util.{HighResolutionClock, Logging} - -/** - * A key value storage engine factory implementation - * - * This trait encapsulates all the steps needed to create a key value storage engine. It is meant to be extended - * by the specific key value store factory implementations which will in turn override the getKVStore method. - */ -trait BaseKeyValueStorageEngineFactory[K, V] extends StorageEngineFactory[K, V] { - - private val INMEMORY_KV_STORAGE_ENGINE_FACTORY = - "org.apache.samza.storage.kv.inmemory.InMemoryKeyValueStorageEngineFactory" - - /** - * Return a KeyValueStore instance for the given store name, - * which will be used as the underlying raw store - * - * @param storeName Name of the store - * @param storeDir The directory of the store - * @param registry MetricsRegistry to which to publish store specific metrics. - * @param changeLogSystemStreamPartition Samza stream partition from which to receive the changelog. - * @param containerContext Information about the container in which the task is executing. - * @return A valid KeyValueStore instance - */ - def getKVStore(storeName: String, - storeDir: File, - registry: MetricsRegistry, - changeLogSystemStreamPartition: SystemStreamPartition, - jobContext: JobContext, - containerContext: ContainerContext, storeMode: StoreMode): KeyValueStore[Array[Byte], Array[Byte]] - - /** - * Constructs a key-value StorageEngine and returns it to the caller - * - * @param storeName The name of the storage engine. - * @param storeDir The directory of the storage engine. - * @param keySerde The serializer to use for serializing keys when reading or writing to the store. - * @param msgSerde The serializer to use for serializing messages when reading or writing to the store. - * @param changelogCollector MessageCollector the storage engine uses to persist changes. - * @param registry MetricsRegistry to which to publish storage-engine specific metrics. - * @param changelogSSP Samza system stream partition from which to receive the changelog. - * @param containerContext Information about the container in which the task is executing. - **/ - def getStorageEngine(storeName: String, - storeDir: File, - keySerde: Serde[K], - msgSerde: Serde[V], - changelogCollector: MessageCollector, - registry: MetricsRegistry, - changelogSSP: SystemStreamPartition, - jobContext: JobContext, - containerContext: ContainerContext, storeMode : StoreMode): StorageEngine = { - val storageConfigSubset = jobContext.getConfig.subset("stores." + storeName + ".", true) - val storageConfig = new StorageConfig(jobContext.getConfig) - val storeFactory = JavaOptionals.toRichOptional(storageConfig.getStorageFactoryClassName(storeName)).toOption - var storePropertiesBuilder = new StoreProperties.StorePropertiesBuilder() - val accessLog = storageConfig.getAccessLogEnabled(storeName) - - var maxMessageSize = storageConfig.getChangelogMaxMsgSizeBytes(storeName) - val disallowLargeMessages = storageConfig.getDisallowLargeMessages(storeName) - val dropLargeMessage = storageConfig.getDropLargeMessages(storeName) - - if (storeFactory.isEmpty) { - throw new SamzaException("Store factory not defined. Cannot proceed with KV store creation!") - } - if (!storeFactory.get.equals(INMEMORY_KV_STORAGE_ENGINE_FACTORY)) { - storePropertiesBuilder = storePropertiesBuilder.setPersistedToDisk(true) - } - - val batchSize = storageConfigSubset.getInt("write.batch.size", 500) - val cacheSize = storageConfigSubset.getInt("object.cache.size", math.max(batchSize, 1000)) - val enableCache = cacheSize > 0 - - if (cacheSize > 0 && cacheSize < batchSize) { - throw new SamzaException("A store's cache.size cannot be less than batch.size as batched values reside in cache.") - } - - if (keySerde == null) { - throw new SamzaException("Must define a key serde when using key value storage.") - } - - if (msgSerde == null) { - throw new SamzaException("Must define a message serde when using key value storage.") - } - - val rawStore = - getKVStore(storeName, storeDir, registry, changelogSSP, jobContext, containerContext, storeMode) - - // maybe wrap with logging - val maybeLoggedStore = if (changelogSSP == null) { - rawStore - } else { - val loggedStoreMetrics = new LoggedStoreMetrics(storeName, registry) - storePropertiesBuilder = storePropertiesBuilder.setLoggedStore(true) - new LoggedStore(rawStore, changelogSSP, changelogCollector, loggedStoreMetrics) - } - - var toBeAccessLoggedStore: KeyValueStore[K, V] = null - - // If large messages are disallowed in config, then this creates a LargeMessageSafeKeyValueStore that throws a - // RecordTooLargeException when a large message is encountered. - if (disallowLargeMessages) { - // maybe wrap with caching - val maybeCachedStore = if (enableCache) { - createCachedStore(storeName, registry, maybeLoggedStore, cacheSize, batchSize) - } else { - maybeLoggedStore - } - - // wrap with large message checking - val largeMessageSafeKeyValueStore = new LargeMessageSafeStore(maybeCachedStore, storeName, false, maxMessageSize) - // wrap with serialization - val serializedMetrics = new SerializedKeyValueStoreMetrics(storeName, registry) - toBeAccessLoggedStore = new SerializedKeyValueStore[K, V](largeMessageSafeKeyValueStore, keySerde, msgSerde, serializedMetrics) - - } - else { - val toBeSerializedStore = if (dropLargeMessage) { - // wrap with large message checking - new LargeMessageSafeStore(maybeLoggedStore, storeName, dropLargeMessage, maxMessageSize) - } else { - maybeLoggedStore - } - // wrap with serialization - val serializedMetrics = new SerializedKeyValueStoreMetrics(storeName, registry) - val serializedStore = new SerializedKeyValueStore[K, V](toBeSerializedStore, keySerde, msgSerde, serializedMetrics) - // maybe wrap with caching - toBeAccessLoggedStore = if (enableCache) { - createCachedStore(storeName, registry, serializedStore, cacheSize, batchSize) - } else { - serializedStore - } - } - - val maybeAccessLoggedStore = if (accessLog) { - new AccessLoggedStore(toBeAccessLoggedStore, changelogCollector, changelogSSP, storageConfig, storeName, keySerde) - } else { - toBeAccessLoggedStore - } - - // wrap with null value checking - val nullSafeStore = new NullSafeKeyValueStore(maybeAccessLoggedStore) - - // create the storage engine and return - val keyValueStorageEngineMetrics = new KeyValueStorageEngineMetrics(storeName, registry) - val metricsConfig = new MetricsConfig(jobContext.getConfig) - val clock = if (metricsConfig.getMetricsTimerEnabled) { - new HighResolutionClock { - override def nanoTime(): Long = System.nanoTime() - } - } else { - new HighResolutionClock { - override def nanoTime(): Long = 0L - } - } - - new KeyValueStorageEngine(storeName, storeDir, storePropertiesBuilder.build(), nullSafeStore, rawStore, - changelogSSP, changelogCollector, keyValueStorageEngineMetrics, batchSize, () => clock.nanoTime()) - } - - def createCachedStore[K, V](storeName: String, registry: MetricsRegistry, - underlyingStore: KeyValueStore[K, V], cacheSize: Int, batchSize: Int): KeyValueStore[K, V] = { - // wrap with caching - val cachedStoreMetrics = new CachedStoreMetrics(storeName, registry) - new CachedStore(underlyingStore, cacheSize, batchSize, cachedStoreMetrics) - } -} diff --git a/samza-kv/src/main/scala/org/apache/samza/storage/kv/CachedStore.scala b/samza-kv/src/main/scala/org/apache/samza/storage/kv/CachedStore.scala index 5c1961cbe0..6dbc21f5ca 100644 --- a/samza-kv/src/main/scala/org/apache/samza/storage/kv/CachedStore.scala +++ b/samza-kv/src/main/scala/org/apache/samza/storage/kv/CachedStore.scala @@ -25,6 +25,7 @@ import scala.collection._ import java.nio.file.Path import java.util.{Arrays, Optional} +import com.google.common.annotations.VisibleForTesting import org.apache.samza.checkpoint.CheckpointId /** @@ -299,6 +300,11 @@ class CachedStore[K, V]( override def checkpoint(id: CheckpointId): Optional[Path] = { store.checkpoint(id) } + + @VisibleForTesting + private[kv] def getStore: KeyValueStore[K, V] = { + store + } } private case class CacheEntry[K, V](var value: V, var dirty: mutable.DoubleLinkedList[K]) diff --git a/samza-kv/src/main/scala/org/apache/samza/storage/kv/KeyValueStorageEngine.scala b/samza-kv/src/main/scala/org/apache/samza/storage/kv/KeyValueStorageEngine.scala index afba824d3a..ddf9ac85e3 100644 --- a/samza-kv/src/main/scala/org/apache/samza/storage/kv/KeyValueStorageEngine.scala +++ b/samza-kv/src/main/scala/org/apache/samza/storage/kv/KeyValueStorageEngine.scala @@ -29,6 +29,7 @@ import org.apache.samza.util.TimerUtil import java.nio.file.Path import java.util.Optional +import com.google.common.annotations.VisibleForTesting import org.apache.samza.checkpoint.CheckpointId /** @@ -252,4 +253,14 @@ class KeyValueStorageEngine[K, V]( wrapperStore.snapshot(from, to) } } + + @VisibleForTesting + private[kv] def getRawStore: KeyValueStore[Array[Byte], Array[Byte]] = { + rawStore + } + + @VisibleForTesting + private[kv] def getWrapperStore: KeyValueStore[K, V] = { + wrapperStore + } } diff --git a/samza-kv/src/main/scala/org/apache/samza/storage/kv/LoggedStore.scala b/samza-kv/src/main/scala/org/apache/samza/storage/kv/LoggedStore.scala index 320e801fed..2db754ad4d 100644 --- a/samza-kv/src/main/scala/org/apache/samza/storage/kv/LoggedStore.scala +++ b/samza-kv/src/main/scala/org/apache/samza/storage/kv/LoggedStore.scala @@ -22,6 +22,7 @@ package org.apache.samza.storage.kv import java.nio.file.Path import java.util.Optional +import com.google.common.annotations.VisibleForTesting import org.apache.samza.checkpoint.CheckpointId import org.apache.samza.util.Logging import org.apache.samza.system.{OutgoingMessageEnvelope, SystemStreamPartition} @@ -125,4 +126,9 @@ class LoggedStore[K, V]( override def checkpoint(id: CheckpointId): Optional[Path] = { store.checkpoint(id) } + + @VisibleForTesting + private[kv] def getStore: KeyValueStore[K, V] = { + store + } } \ No newline at end of file diff --git a/samza-kv/src/main/scala/org/apache/samza/storage/kv/NullSafeKeyValueStore.scala b/samza-kv/src/main/scala/org/apache/samza/storage/kv/NullSafeKeyValueStore.scala index 8bb6fa2a1c..feca04b620 100644 --- a/samza-kv/src/main/scala/org/apache/samza/storage/kv/NullSafeKeyValueStore.scala +++ b/samza-kv/src/main/scala/org/apache/samza/storage/kv/NullSafeKeyValueStore.scala @@ -22,6 +22,7 @@ package org.apache.samza.storage.kv import java.nio.file.Path import java.util.Optional +import com.google.common.annotations.VisibleForTesting import org.apache.samza.checkpoint.CheckpointId import scala.collection.JavaConverters._ @@ -104,4 +105,9 @@ class NullSafeKeyValueStore[K, V](store: KeyValueStore[K, V]) extends KeyValueSt override def checkpoint(id: CheckpointId): Optional[Path] = { store.checkpoint(id) } + + @VisibleForTesting + private[kv] def getStore: KeyValueStore[K, V] = { + store + } } diff --git a/samza-kv/src/main/scala/org/apache/samza/storage/kv/SerializedKeyValueStore.scala b/samza-kv/src/main/scala/org/apache/samza/storage/kv/SerializedKeyValueStore.scala index 96566ac4e7..b78e14a88a 100644 --- a/samza-kv/src/main/scala/org/apache/samza/storage/kv/SerializedKeyValueStore.scala +++ b/samza-kv/src/main/scala/org/apache/samza/storage/kv/SerializedKeyValueStore.scala @@ -22,6 +22,7 @@ package org.apache.samza.storage.kv import java.nio.file.Path import java.util.Optional +import com.google.common.annotations.VisibleForTesting import org.apache.samza.checkpoint.CheckpointId import org.apache.samza.util.Logging import org.apache.samza.serializers._ @@ -202,4 +203,9 @@ class SerializedKeyValueStore[K, V]( override def checkpoint(id: CheckpointId): Optional[Path] = { store.checkpoint(id) } + + @VisibleForTesting + private[kv] def getStore: KeyValueStore[Array[Byte], Array[Byte]] = { + store + } } diff --git a/samza-kv/src/test/java/org/apache/samza/storage/kv/MockKeyValueStorageEngineFactory.java b/samza-kv/src/test/java/org/apache/samza/storage/kv/MockKeyValueStorageEngineFactory.java new file mode 100644 index 0000000000..3430ae951e --- /dev/null +++ b/samza-kv/src/test/java/org/apache/samza/storage/kv/MockKeyValueStorageEngineFactory.java @@ -0,0 +1,45 @@ +/* + * 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.samza.storage.kv; + +import java.io.File; +import org.apache.samza.context.ContainerContext; +import org.apache.samza.context.JobContext; +import org.apache.samza.metrics.MetricsRegistry; +import org.apache.samza.system.SystemStreamPartition; + + +/** + * Used for testing {@link BaseKeyValueStorageEngineFactory}. + * Implements {@link #getKVStore} to return a pre-built {@link KeyValueStore}. + */ +public class MockKeyValueStorageEngineFactory extends BaseKeyValueStorageEngineFactory { + private final KeyValueStore rawKeyValueStore; + + public MockKeyValueStorageEngineFactory(KeyValueStore rawKeyValueStore) { + this.rawKeyValueStore = rawKeyValueStore; + } + + @Override + protected KeyValueStore getKVStore(String storeName, File storeDir, MetricsRegistry registry, + SystemStreamPartition changeLogSystemStreamPartition, JobContext jobContext, ContainerContext containerContext, + StoreMode storeMode) { + return this.rawKeyValueStore; + } +} diff --git a/samza-kv/src/test/java/org/apache/samza/storage/kv/TestBaseKeyValueStorageEngineFactory.java b/samza-kv/src/test/java/org/apache/samza/storage/kv/TestBaseKeyValueStorageEngineFactory.java new file mode 100644 index 0000000000..46eab03669 --- /dev/null +++ b/samza-kv/src/test/java/org/apache/samza/storage/kv/TestBaseKeyValueStorageEngineFactory.java @@ -0,0 +1,273 @@ +/* + * 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.samza.storage.kv; + +import java.io.File; +import java.util.Map; +import com.google.common.collect.ImmutableMap; +import org.apache.samza.Partition; +import org.apache.samza.SamzaException; +import org.apache.samza.config.Config; +import org.apache.samza.config.MapConfig; +import org.apache.samza.config.StorageConfig; +import org.apache.samza.context.ContainerContext; +import org.apache.samza.context.JobContext; +import org.apache.samza.metrics.Gauge; +import org.apache.samza.metrics.MetricsRegistry; +import org.apache.samza.serializers.Serde; +import org.apache.samza.storage.StorageEngine; +import org.apache.samza.storage.StorageEngineFactory; +import org.apache.samza.storage.StoreProperties; +import org.apache.samza.system.SystemStreamPartition; +import org.apache.samza.task.MessageCollector; +import org.junit.Before; +import org.junit.Test; +import org.mockito.Mock; +import org.mockito.MockitoAnnotations; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; +import static org.mockito.Matchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + + +public class TestBaseKeyValueStorageEngineFactory { + private static final String STORE_NAME = "myStore"; + private static final StorageEngineFactory.StoreMode STORE_MODE = StorageEngineFactory.StoreMode.ReadWrite; + private static final SystemStreamPartition CHANGELOG_SSP = + new SystemStreamPartition("system", "stream", new Partition(0)); + private static final Map BASE_CONFIG = + ImmutableMap.of(String.format(StorageConfig.FACTORY, STORE_NAME), + MockKeyValueStorageEngineFactory.class.getName()); + private static final Map DISABLE_CACHE = + ImmutableMap.of(String.format("stores.%s.object.cache.size", STORE_NAME), "0"); + private static final Map DISALLOW_LARGE_MESSAGES = + ImmutableMap.of(String.format(StorageConfig.DISALLOW_LARGE_MESSAGES, STORE_NAME), "true"); + private static final Map DROP_LARGE_MESSAGES = + ImmutableMap.of(String.format(StorageConfig.DROP_LARGE_MESSAGES, STORE_NAME), "true"); + private static final Map ACCESS_LOG_ENABLED = + ImmutableMap.of(String.format("stores.%s.accesslog.enabled", STORE_NAME), "true"); + + @Mock + private File storeDir; + @Mock + private Serde keySerde; + @Mock + private Serde msgSerde; + @Mock + private MessageCollector changelogCollector; + @Mock + private MetricsRegistry metricsRegistry; + @Mock + private JobContext jobContext; + @Mock + private ContainerContext containerContext; + @Mock + private KeyValueStore rawKeyValueStore; + + @Before + public void setup() { + MockitoAnnotations.initMocks(this); + // some metrics objects need this for histogram metric instantiation + when(this.metricsRegistry.newGauge(any(), any())).thenReturn(mock(Gauge.class)); + } + + @Test(expected = SamzaException.class) + public void testMissingStoreFactory() { + Config config = new MapConfig(); + callGetStorageEngine(config, null); + } + + @Test(expected = SamzaException.class) + public void testInvalidCacheSize() { + Config config = new MapConfig(BASE_CONFIG, + ImmutableMap.of(String.format("stores.%s.write.cache.batch", STORE_NAME), "100", + String.format("stores.%s.object.cache.size", STORE_NAME), "50")); + callGetStorageEngine(config, null); + } + + @Test + public void testInMemoryKeyValueStore() { + Config config = new MapConfig(DISABLE_CACHE, ImmutableMap.of(String.format(StorageConfig.FACTORY, STORE_NAME), + "org.apache.samza.storage.kv.inmemory.InMemoryKeyValueStorageEngineFactory")); + StorageEngine storageEngine = callGetStorageEngine(config, null); + KeyValueStorageEngine keyValueStorageEngine = baseStorageEngineValidation(storageEngine); + assertStoreProperties(keyValueStorageEngine.getStoreProperties(), false, false); + NullSafeKeyValueStore nullSafeKeyValueStore = + assertAndCast(keyValueStorageEngine.getWrapperStore(), NullSafeKeyValueStore.class); + SerializedKeyValueStore serializedKeyValueStore = + assertAndCast(nullSafeKeyValueStore.getStore(), SerializedKeyValueStore.class); + // config has the in-memory key-value factory, but still calling the test factory, so store will be the test store + assertEquals(this.rawKeyValueStore, serializedKeyValueStore.getStore()); + } + + @Test + public void testRawStoreOnly() { + Config config = new MapConfig(BASE_CONFIG, DISABLE_CACHE); + StorageEngine storageEngine = callGetStorageEngine(config, null); + KeyValueStorageEngine keyValueStorageEngine = baseStorageEngineValidation(storageEngine); + assertStoreProperties(keyValueStorageEngine.getStoreProperties(), true, false); + NullSafeKeyValueStore nullSafeKeyValueStore = + assertAndCast(keyValueStorageEngine.getWrapperStore(), NullSafeKeyValueStore.class); + SerializedKeyValueStore serializedKeyValueStore = + assertAndCast(nullSafeKeyValueStore.getStore(), SerializedKeyValueStore.class); + assertEquals(this.rawKeyValueStore, serializedKeyValueStore.getStore()); + } + + @Test + public void testWithLoggedStore() { + Config config = new MapConfig(BASE_CONFIG, DISABLE_CACHE); + StorageEngine storageEngine = callGetStorageEngine(config, CHANGELOG_SSP); + KeyValueStorageEngine keyValueStorageEngine = baseStorageEngineValidation(storageEngine); + assertStoreProperties(keyValueStorageEngine.getStoreProperties(), true, true); + NullSafeKeyValueStore nullSafeKeyValueStore = + assertAndCast(keyValueStorageEngine.getWrapperStore(), NullSafeKeyValueStore.class); + SerializedKeyValueStore serializedKeyValueStore = + assertAndCast(nullSafeKeyValueStore.getStore(), SerializedKeyValueStore.class); + LoggedStore loggedStore = assertAndCast(serializedKeyValueStore.getStore(), LoggedStore.class); + // noinspection AssertEqualsBetweenInconvertibleTypes + assertEquals(this.rawKeyValueStore, loggedStore.getStore()); + } + + @Test + public void testWithLoggedStoreWithCache() { + Config config = new MapConfig(BASE_CONFIG); + StorageEngine storageEngine = callGetStorageEngine(config, CHANGELOG_SSP); + KeyValueStorageEngine keyValueStorageEngine = baseStorageEngineValidation(storageEngine); + assertStoreProperties(keyValueStorageEngine.getStoreProperties(), true, true); + NullSafeKeyValueStore nullSafeKeyValueStore = + assertAndCast(keyValueStorageEngine.getWrapperStore(), NullSafeKeyValueStore.class); + CachedStore cachedStore = assertAndCast(nullSafeKeyValueStore.getStore(), CachedStore.class); + SerializedKeyValueStore serializedKeyValueStore = + assertAndCast(cachedStore.getStore(), SerializedKeyValueStore.class); + LoggedStore loggedStore = assertAndCast(serializedKeyValueStore.getStore(), LoggedStore.class); + // noinspection AssertEqualsBetweenInconvertibleTypes + assertEquals(this.rawKeyValueStore, loggedStore.getStore()); + } + + @Test + public void testDisallowLargeMessages() { + Config config = new MapConfig(BASE_CONFIG, DISABLE_CACHE, DISALLOW_LARGE_MESSAGES); + StorageEngine storageEngine = callGetStorageEngine(config, null); + KeyValueStorageEngine keyValueStorageEngine = baseStorageEngineValidation(storageEngine); + assertStoreProperties(keyValueStorageEngine.getStoreProperties(), true, false); + NullSafeKeyValueStore nullSafeKeyValueStore = + assertAndCast(keyValueStorageEngine.getWrapperStore(), NullSafeKeyValueStore.class); + SerializedKeyValueStore serializedKeyValueStore = + assertAndCast(nullSafeKeyValueStore.getStore(), SerializedKeyValueStore.class); + LargeMessageSafeStore largeMessageSafeStore = + assertAndCast(serializedKeyValueStore.getStore(), LargeMessageSafeStore.class); + assertEquals(this.rawKeyValueStore, largeMessageSafeStore.getStore()); + } + + @Test + public void testDisallowLargeMessagesWithCache() { + Config config = new MapConfig(BASE_CONFIG, DISALLOW_LARGE_MESSAGES); + StorageEngine storageEngine = callGetStorageEngine(config, null); + KeyValueStorageEngine keyValueStorageEngine = baseStorageEngineValidation(storageEngine); + assertStoreProperties(keyValueStorageEngine.getStoreProperties(), true, false); + NullSafeKeyValueStore nullSafeKeyValueStore = + assertAndCast(keyValueStorageEngine.getWrapperStore(), NullSafeKeyValueStore.class); + SerializedKeyValueStore serializedKeyValueStore = + assertAndCast(nullSafeKeyValueStore.getStore(), SerializedKeyValueStore.class); + LargeMessageSafeStore largeMessageSafeStore = + assertAndCast(serializedKeyValueStore.getStore(), LargeMessageSafeStore.class); + CachedStore cachedStore = assertAndCast(largeMessageSafeStore.getStore(), CachedStore.class); + // noinspection AssertEqualsBetweenInconvertibleTypes + assertEquals(this.rawKeyValueStore, cachedStore.getStore()); + } + + @Test + public void testDropLargeMessages() { + Config config = new MapConfig(BASE_CONFIG, DISABLE_CACHE, DROP_LARGE_MESSAGES); + StorageEngine storageEngine = callGetStorageEngine(config, null); + KeyValueStorageEngine keyValueStorageEngine = baseStorageEngineValidation(storageEngine); + assertStoreProperties(keyValueStorageEngine.getStoreProperties(), true, false); + NullSafeKeyValueStore nullSafeKeyValueStore = + assertAndCast(keyValueStorageEngine.getWrapperStore(), NullSafeKeyValueStore.class); + SerializedKeyValueStore serializedKeyValueStore = + assertAndCast(nullSafeKeyValueStore.getStore(), SerializedKeyValueStore.class); + LargeMessageSafeStore largeMessageSafeStore = + assertAndCast(serializedKeyValueStore.getStore(), LargeMessageSafeStore.class); + assertEquals(this.rawKeyValueStore, largeMessageSafeStore.getStore()); + } + + @Test + public void testDropLargeMessagesWithCache() { + Config config = new MapConfig(BASE_CONFIG, DROP_LARGE_MESSAGES); + StorageEngine storageEngine = callGetStorageEngine(config, null); + KeyValueStorageEngine keyValueStorageEngine = baseStorageEngineValidation(storageEngine); + assertStoreProperties(keyValueStorageEngine.getStoreProperties(), true, false); + NullSafeKeyValueStore nullSafeKeyValueStore = + assertAndCast(keyValueStorageEngine.getWrapperStore(), NullSafeKeyValueStore.class); + CachedStore cachedStore = assertAndCast(nullSafeKeyValueStore.getStore(), CachedStore.class); + SerializedKeyValueStore serializedKeyValueStore = + assertAndCast(cachedStore.getStore(), SerializedKeyValueStore.class); + LargeMessageSafeStore largeMessageSafeStore = + assertAndCast(serializedKeyValueStore.getStore(), LargeMessageSafeStore.class); + assertEquals(this.rawKeyValueStore, largeMessageSafeStore.getStore()); + } + + @Test + public void testAccessLogStore() { + Config config = new MapConfig(BASE_CONFIG, DISABLE_CACHE, ACCESS_LOG_ENABLED); + // AccessLoggedStore requires a changelog SSP + StorageEngine storageEngine = callGetStorageEngine(config, CHANGELOG_SSP); + KeyValueStorageEngine keyValueStorageEngine = baseStorageEngineValidation(storageEngine); + assertStoreProperties(keyValueStorageEngine.getStoreProperties(), true, true); + NullSafeKeyValueStore nullSafeKeyValueStore = + assertAndCast(keyValueStorageEngine.getWrapperStore(), NullSafeKeyValueStore.class); + AccessLoggedStore accessLoggedStore = + assertAndCast(nullSafeKeyValueStore.getStore(), AccessLoggedStore.class); + SerializedKeyValueStore serializedKeyValueStore = + assertAndCast(accessLoggedStore.getStore(), SerializedKeyValueStore.class); + LoggedStore loggedStore = assertAndCast(serializedKeyValueStore.getStore(), LoggedStore.class); + // noinspection AssertEqualsBetweenInconvertibleTypes + assertEquals(this.rawKeyValueStore, loggedStore.getStore()); + } + + private static > T assertAndCast(KeyValueStore keyValueStore, Class clazz) { + assertTrue("Expected type " + clazz.getName(), clazz.isInstance(keyValueStore)); + return clazz.cast(keyValueStore); + } + + private KeyValueStorageEngine baseStorageEngineValidation(StorageEngine storageEngine) { + assertTrue(storageEngine instanceof KeyValueStorageEngine); + KeyValueStorageEngine keyValueStorageEngine = (KeyValueStorageEngine) storageEngine; + assertEquals(this.rawKeyValueStore, keyValueStorageEngine.getRawStore()); + return keyValueStorageEngine; + } + + private static void assertStoreProperties(StoreProperties storeProperties, boolean expectedPersistedToDisk, + boolean expectedLoggedStore) { + assertEquals(expectedPersistedToDisk, storeProperties.isPersistedToDisk()); + assertEquals(expectedLoggedStore, storeProperties.isLoggedStore()); + } + + /** + * @param changelogSSP if non-null, then enables logged store + */ + private StorageEngine callGetStorageEngine(Config config, SystemStreamPartition changelogSSP) { + when(this.jobContext.getConfig()).thenReturn(config); + return new MockKeyValueStorageEngineFactory(this.rawKeyValueStore).getStorageEngine(STORE_NAME, this.storeDir, + this.keySerde, this.msgSerde, this.changelogCollector, this.metricsRegistry, changelogSSP, this.jobContext, + this.containerContext, STORE_MODE); + } +} From 76859a95b2eb566d88b14d4023521d90ecc6588d Mon Sep 17 00:00:00 2001 From: Cameron Lee Date: Wed, 27 May 2020 18:21:23 -0700 Subject: [PATCH 2/2] added javadocs, extract constants, improve exception messages --- .../kv/BaseKeyValueStorageEngineFactory.java | 61 ++++++++++++++++--- .../TestBaseKeyValueStorageEngineFactory.java | 24 +++++++- 2 files changed, 74 insertions(+), 11 deletions(-) diff --git a/samza-kv/src/main/java/org/apache/samza/storage/kv/BaseKeyValueStorageEngineFactory.java b/samza-kv/src/main/java/org/apache/samza/storage/kv/BaseKeyValueStorageEngineFactory.java index 44ada3cede..a3fc1783bf 100644 --- a/samza-kv/src/main/java/org/apache/samza/storage/kv/BaseKeyValueStorageEngineFactory.java +++ b/samza-kv/src/main/java/org/apache/samza/storage/kv/BaseKeyValueStorageEngineFactory.java @@ -46,6 +46,10 @@ public abstract class BaseKeyValueStorageEngineFactory implements StorageEngineFactory { private static final String INMEMORY_KV_STORAGE_ENGINE_FACTORY = "org.apache.samza.storage.kv.inmemory.InMemoryKeyValueStorageEngineFactory"; + private static final String WRITE_BATCH_SIZE = "write.batch.size"; + private static final int DEFAULT_WRITE_BATCH_SIZE = 500; + private static final String OBJECT_CACHE_SIZE = "object.cache.size"; + private static final int DEFAULT_OBJECT_CACHE_SIZE = 1000; /** * Implement this to return a KeyValueStore instance for the given store name, which will be used as the underlying @@ -94,22 +98,26 @@ public StorageEngine getStorageEngine(String storeName, Optional storeFactory = storageConfig.getStorageFactoryClassName(storeName); StoreProperties.StorePropertiesBuilder storePropertiesBuilder = new StoreProperties.StorePropertiesBuilder(); if (!storeFactory.isPresent() || StringUtils.isBlank(storeFactory.get())) { - throw new SamzaException("Store factory not defined. Cannot proceed with KV store creation!"); + throw new SamzaException( + String.format("Store factory not defined for store %s. Cannot proceed with KV store creation!", storeName)); } if (!storeFactory.get().equals(INMEMORY_KV_STORAGE_ENGINE_FACTORY)) { storePropertiesBuilder.setPersistedToDisk(true); } - int batchSize = storageConfigSubset.getInt("write.batch.size", 500); - int cacheSize = storageConfigSubset.getInt("object.cache.size", Math.max(batchSize, 1000)); + int batchSize = storageConfigSubset.getInt(WRITE_BATCH_SIZE, DEFAULT_WRITE_BATCH_SIZE); + int cacheSize = storageConfigSubset.getInt(OBJECT_CACHE_SIZE, Math.max(batchSize, DEFAULT_OBJECT_CACHE_SIZE)); if (cacheSize > 0 && cacheSize < batchSize) { throw new SamzaException( - "A store's cache.size cannot be less than batch.size as batched values reside in cache."); + String.format("cache.size for store %s cannot be less than batch.size as batched values reside in cache.", + storeName)); } if (keySerde == null) { - throw new SamzaException("Must define a key serde when using key value storage."); + throw new SamzaException( + String.format("Must define a key serde when using key value storage for store %s.", storeName)); } if (msgSerde == null) { - throw new SamzaException("Must define a message serde when using key value storage."); + throw new SamzaException( + String.format("Must define a message serde when using key value storage for store %s.", storeName)); } KeyValueStore rawStore = @@ -117,7 +125,7 @@ public StorageEngine getStorageEngine(String storeName, KeyValueStore maybeLoggedStore = buildMaybeLoggedStore(changelogSSP, storeName, registry, storePropertiesBuilder, rawStore, changelogCollector); // this also applies serialization and caching layers - KeyValueStore toBeAccessLoggedStore = applyLargeMessageHandling(storeName, registry, + KeyValueStore toBeAccessLoggedStore = buildStoreWithLargeMessageHandling(storeName, registry, maybeLoggedStore, storageConfig, cacheSize, batchSize, keySerde, msgSerde); KeyValueStore maybeAccessLoggedStore = buildMaybeAccessLoggedStore(storeName, toBeAccessLoggedStore, changelogCollector, changelogSSP, storageConfig, @@ -131,6 +139,10 @@ public StorageEngine getStorageEngine(String storeName, ScalaJavaUtil.toScalaFunction(clock::nanoTime)); } + /** + * Wraps {@code storeToWrap} into a {@link LoggedStore} if {@code changelogSSP} is defined. + * Otherwise, returns the original {@code storeToWrap}. + */ private static KeyValueStore buildMaybeLoggedStore(SystemStreamPartition changelogSSP, String storeName, MetricsRegistry registry, @@ -146,7 +158,14 @@ private static KeyValueStore buildMaybeLoggedStore(SystemStreamP } } - private static KeyValueStore applyLargeMessageHandling(String storeName, + /** + * Wraps {@code storeToWrap} with the proper layers to handle large messages. + * If "disallow.large.messages" is enabled, then the message will be serialized and the size will be checked before + * storing in the serialized message in the cache. + * If "disallow.large.messages" is disabled, then the deserialized message will be stored in the cache. If + * "drop.large.messages" is enabled, then large messages will not be sent to the logged store. + */ + private static KeyValueStore buildStoreWithLargeMessageHandling(String storeName, MetricsRegistry registry, KeyValueStore storeToWrap, StorageConfig storageConfig, @@ -157,11 +176,13 @@ private static KeyValueStore applyLargeMessageHandling(String store int maxMessageSize = storageConfig.getChangelogMaxMsgSizeBytes(storeName); if (storageConfig.getDisallowLargeMessages(storeName)) { /* - * If large messages are disallowed in config, then this creates a LargeMessageSafeKeyValueStore that throws a - * RecordTooLargeException when a large message is encountered. + * The store wrapping ordering is done this way so that a large message cannot end up in the cache. However, it + * also means that serialized data is in the cache, so performance will be worse since the data needs to be + * deserialized even when cached. */ KeyValueStore maybeCachedStore = buildMaybeCachedStore(storeName, registry, storeToWrap, cacheSize, batchSize); + // this will throw a RecordTooLargeException when a large message is encountered LargeMessageSafeStore largeMessageSafeKeyValueStore = new LargeMessageSafeStore(maybeCachedStore, storeName, false, maxMessageSize); return buildSerializedStore(storeName, registry, largeMessageSafeKeyValueStore, keySerde, msgSerde); @@ -174,10 +195,18 @@ private static KeyValueStore applyLargeMessageHandling(String store } KeyValueStore serializedStore = buildSerializedStore(storeName, registry, toBeSerializedStore, keySerde, msgSerde); + /* + * Allows deserialized entries to be stored in the cache, but it means that a large message may end up in the + * cache even though it was not persisted to the logged store. + */ return buildMaybeCachedStore(storeName, registry, serializedStore, cacheSize, batchSize); } } + /** + * Wraps {@code storeToWrap} with a {@link CachedStore} if caching is enabled. + * Otherwise, returns the {@code storeToWrap}. + */ private static KeyValueStore buildMaybeCachedStore(String storeName, MetricsRegistry registry, KeyValueStore storeToWrap, int cacheSize, int batchSize) { if (cacheSize > 0) { @@ -188,6 +217,9 @@ private static KeyValueStore buildMaybeCachedStore(String storeName } } + /** + * Wraps {@code storeToWrap} with a {@link SerializedKeyValueStore}. + */ private static KeyValueStore buildSerializedStore(String storeName, MetricsRegistry registry, KeyValueStore storeToWrap, @@ -197,6 +229,10 @@ private static KeyValueStore buildSerializedStore(String storeName, return new SerializedKeyValueStore<>(storeToWrap, keySerde, msgSerde, serializedMetrics); } + /** + * Wraps {@code storeToWrap} with an {@link AccessLoggedStore} if enabled. + * Otherwise, returns the {@code storeToWrap}. + */ private static KeyValueStore buildMaybeAccessLoggedStore(String storeName, KeyValueStore storeToWrap, MessageCollector changelogCollector, @@ -210,6 +246,11 @@ private static KeyValueStore buildMaybeAccessLoggedStore(String sto } } + /** + * If "metrics.timer.enabled" is enabled, then returns a {@link HighResolutionClock} that uses + * {@link System#nanoTime}. + * Otherwise, returns a clock which always returns 0. + */ private static HighResolutionClock buildClock(Config config) { MetricsConfig metricsConfig = new MetricsConfig(config); if (metricsConfig.getMetricsTimerEnabled()) { diff --git a/samza-kv/src/test/java/org/apache/samza/storage/kv/TestBaseKeyValueStorageEngineFactory.java b/samza-kv/src/test/java/org/apache/samza/storage/kv/TestBaseKeyValueStorageEngineFactory.java index 46eab03669..22a1b57027 100644 --- a/samza-kv/src/test/java/org/apache/samza/storage/kv/TestBaseKeyValueStorageEngineFactory.java +++ b/samza-kv/src/test/java/org/apache/samza/storage/kv/TestBaseKeyValueStorageEngineFactory.java @@ -103,6 +103,24 @@ public void testInvalidCacheSize() { callGetStorageEngine(config, null); } + @Test(expected = SamzaException.class) + public void testMissingKeySerde() { + Config config = new MapConfig(BASE_CONFIG); + when(this.jobContext.getConfig()).thenReturn(config); + new MockKeyValueStorageEngineFactory(this.rawKeyValueStore).getStorageEngine(STORE_NAME, this.storeDir, null, + this.msgSerde, this.changelogCollector, this.metricsRegistry, null, this.jobContext, this.containerContext, + STORE_MODE); + } + + @Test(expected = SamzaException.class) + public void testMissingValueSerde() { + Config config = new MapConfig(BASE_CONFIG); + when(this.jobContext.getConfig()).thenReturn(config); + new MockKeyValueStorageEngineFactory(this.rawKeyValueStore).getStorageEngine(STORE_NAME, this.storeDir, + this.keySerde, null, this.changelogCollector, this.metricsRegistry, null, this.jobContext, + this.containerContext, STORE_MODE); + } + @Test public void testInMemoryKeyValueStore() { Config config = new MapConfig(DISABLE_CACHE, ImmutableMap.of(String.format(StorageConfig.FACTORY, STORE_NAME), @@ -142,12 +160,13 @@ public void testWithLoggedStore() { SerializedKeyValueStore serializedKeyValueStore = assertAndCast(nullSafeKeyValueStore.getStore(), SerializedKeyValueStore.class); LoggedStore loggedStore = assertAndCast(serializedKeyValueStore.getStore(), LoggedStore.class); + // type generics don't match due to wildcard type, but checking reference equality, so type generics don't matter // noinspection AssertEqualsBetweenInconvertibleTypes assertEquals(this.rawKeyValueStore, loggedStore.getStore()); } @Test - public void testWithLoggedStoreWithCache() { + public void testWithLoggedStoreAndCachedStore() { Config config = new MapConfig(BASE_CONFIG); StorageEngine storageEngine = callGetStorageEngine(config, CHANGELOG_SSP); KeyValueStorageEngine keyValueStorageEngine = baseStorageEngineValidation(storageEngine); @@ -158,6 +177,7 @@ public void testWithLoggedStoreWithCache() { SerializedKeyValueStore serializedKeyValueStore = assertAndCast(cachedStore.getStore(), SerializedKeyValueStore.class); LoggedStore loggedStore = assertAndCast(serializedKeyValueStore.getStore(), LoggedStore.class); + // type generics don't match due to wildcard type, but checking reference equality, so type generics don't matter // noinspection AssertEqualsBetweenInconvertibleTypes assertEquals(this.rawKeyValueStore, loggedStore.getStore()); } @@ -190,6 +210,7 @@ public void testDisallowLargeMessagesWithCache() { LargeMessageSafeStore largeMessageSafeStore = assertAndCast(serializedKeyValueStore.getStore(), LargeMessageSafeStore.class); CachedStore cachedStore = assertAndCast(largeMessageSafeStore.getStore(), CachedStore.class); + // type generics don't match due to wildcard type, but checking reference equality, so type generics don't matter // noinspection AssertEqualsBetweenInconvertibleTypes assertEquals(this.rawKeyValueStore, cachedStore.getStore()); } @@ -239,6 +260,7 @@ public void testAccessLogStore() { SerializedKeyValueStore serializedKeyValueStore = assertAndCast(accessLoggedStore.getStore(), SerializedKeyValueStore.class); LoggedStore loggedStore = assertAndCast(serializedKeyValueStore.getStore(), LoggedStore.class); + // type generics don't match due to wildcard type, but checking reference equality, so type generics don't matter // noinspection AssertEqualsBetweenInconvertibleTypes assertEquals(this.rawKeyValueStore, loggedStore.getStore()); }