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) } 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..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 @@ -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,33 @@ 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)) + put(CoordinatorStreamStore.serializeCoordinatorMessageKeyToJson(SetConfig.TYPE, keyForNullVal), + valueSerde.toBytes(nullVal)) + } + + 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)) + } }