From a833a3c9dec73abe77cbf69b53a26fe39f75e414 Mon Sep 17 00:00:00 2001 From: Dengpan Yin Date: Tue, 10 Sep 2019 17:49:46 -0700 Subject: [PATCH 1/3] SAMZA-2318: Empty config values from coordinator stream shouldn't be removed In the previous PR: https://github.com/apache/samza/pull/1010/files config with null/empty values will be filter out when reading the coordinator stream, this will cause some problems when the user expects an empty string to be a value config. --- .../scala/org/apache/samza/util/CoordinatorStreamUtil.scala | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/samza-core/src/main/scala/org/apache/samza/util/CoordinatorStreamUtil.scala b/samza-core/src/main/scala/org/apache/samza/util/CoordinatorStreamUtil.scala index a883b0e838..e38359db9f 100644 --- a/samza-core/src/main/scala/org/apache/samza/util/CoordinatorStreamUtil.scala +++ b/samza-core/src/main/scala/org/apache/samza/util/CoordinatorStreamUtil.scala @@ -131,8 +131,8 @@ object CoordinatorStreamUtil extends Logging { } else { val valueSerde: CoordinatorStreamValueSerde = new CoordinatorStreamValueSerde(SetConfig.TYPE) val valueAsString: String = valueSerde.fromBytes(valueAsBytes) - if (StringUtils.isBlank(valueAsString)) { - warn("Value for key: %s in config is empty or null. Ignoring it." format key) + if (valueAsString == null) { + warn("Value for key: %s in config is decoded to be null. Ignoring it." format key) } else { configMap.put(key, valueAsString) } From 5efb6fe47d92776c9d94cdcd29a8cb363595057d Mon Sep 17 00:00:00 2001 From: Dengpan Yin Date: Thu, 12 Sep 2019 09:33:08 -0700 Subject: [PATCH 2/3] add unit-test --- .../util/TestCoordinatorStreamUtil.scala | 34 ++++++++++++++++++- 1 file changed, 33 insertions(+), 1 deletion(-) diff --git a/samza-core/src/test/scala/org/apache/samza/util/TestCoordinatorStreamUtil.scala b/samza-core/src/test/scala/org/apache/samza/util/TestCoordinatorStreamUtil.scala index 76b077a059..8f207008ff 100644 --- a/samza-core/src/test/scala/org/apache/samza/util/TestCoordinatorStreamUtil.scala +++ b/samza-core/src/test/scala/org/apache/samza/util/TestCoordinatorStreamUtil.scala @@ -18,8 +18,13 @@ */ package org.apache.samza.util +import java.util + +import org.apache.samza.coordinator.metadatastore.CoordinatorStreamStore +import org.apache.samza.coordinator.stream.CoordinatorStreamValueSerde +import org.apache.samza.coordinator.stream.messages.SetConfig import org.apache.samza.system.{StreamSpec, SystemAdmin, SystemStream} -import org.junit.Test +import org.junit.{Assert, Test} import org.mockito.Matchers.any import org.mockito.Mockito @@ -34,4 +39,31 @@ class TestCoordinatorStreamUtil { Mockito.verify(systemStream).getStream Mockito.verify(systemAdmin).createStream(any(classOf[StreamSpec])) } + + @Test + def testReadConfigFromCoordinatorStream { + val keyForNonBlankVal = "app.id" + val nonBlankVal = "1" + val keyForEmptyVal = "task.opt" + val emptyVal = "" + val keyForNullVal = "zk.server" + val nullVal = null + + val valueSerde = new CoordinatorStreamValueSerde(SetConfig.TYPE) + val configMap = new util.HashMap[String, Array[Byte]]() { + put(CoordinatorStreamStore.serializeCoordinatorMessageKeyToJson(SetConfig.TYPE, keyForNonBlankVal), + valueSerde.toBytes(nonBlankVal)) + put(CoordinatorStreamStore.serializeCoordinatorMessageKeyToJson(SetConfig.TYPE, keyForEmptyVal), + valueSerde.toBytes(emptyVal)) + } + + val coordinatorStreamStore = Mockito.mock(classOf[CoordinatorStreamStore]) + Mockito.when(coordinatorStreamStore.all()).thenReturn(configMap) + + val configFromCoordinatorStream = CoordinatorStreamUtil.readConfigFromCoordinatorStream(coordinatorStreamStore) + + Assert.assertEquals(configFromCoordinatorStream.get(keyForNonBlankVal), nonBlankVal) + Assert.assertEquals(configFromCoordinatorStream.get(keyForEmptyVal), emptyVal) + Assert.assertFalse(configFromCoordinatorStream.containsKey(keyForNullVal)) + } } From ec431224767112c9a98d704b378207802a49eae9 Mon Sep 17 00:00:00 2001 From: Dengpan Yin Date: Fri, 13 Sep 2019 10:37:28 -0700 Subject: [PATCH 3/3] Add null value to config map for testing --- .../scala/org/apache/samza/util/TestCoordinatorStreamUtil.scala | 2 ++ 1 file changed, 2 insertions(+) diff --git a/samza-core/src/test/scala/org/apache/samza/util/TestCoordinatorStreamUtil.scala b/samza-core/src/test/scala/org/apache/samza/util/TestCoordinatorStreamUtil.scala index 8f207008ff..eb035960d9 100644 --- a/samza-core/src/test/scala/org/apache/samza/util/TestCoordinatorStreamUtil.scala +++ b/samza-core/src/test/scala/org/apache/samza/util/TestCoordinatorStreamUtil.scala @@ -55,6 +55,8 @@ class TestCoordinatorStreamUtil { valueSerde.toBytes(nonBlankVal)) put(CoordinatorStreamStore.serializeCoordinatorMessageKeyToJson(SetConfig.TYPE, keyForEmptyVal), valueSerde.toBytes(emptyVal)) + put(CoordinatorStreamStore.serializeCoordinatorMessageKeyToJson(SetConfig.TYPE, keyForNullVal), + valueSerde.toBytes(nullVal)) } val coordinatorStreamStore = Mockito.mock(classOf[CoordinatorStreamStore])