Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Might be useful to add a unit test to test that null values appear as non-entries and empty string values are entries with an empty string when loaded from the coordinator stream.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@dnishimura UNIT-Test added Thanks

} else {
configMap.put(key, valueAsString)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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
Comment thread
dengpanyin marked this conversation as resolved.

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))
}
}