Skip to content
Merged
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 @@ -22,6 +22,7 @@
package org.apache.samza.util

import java.util
import org.apache.commons.lang3.StringUtils
import org.apache.samza.SamzaException
import org.apache.samza.config._
import org.apache.samza.system.{SystemFactory, SystemStream}
Expand All @@ -34,7 +35,7 @@ import org.apache.samza.util.ScalaJavaUtil.JavaOptionals
import scala.collection.immutable.Map
import scala.collection.JavaConverters._

object CoordinatorStreamUtil {
object CoordinatorStreamUtil extends Logging {
/**
* Given a job's full config object, build a subset config which includes
* only the job name, job id, and system config for the coordinator stream.
Expand Down Expand Up @@ -106,9 +107,17 @@ object CoordinatorStreamUtil {
val configFromCoordinatorStream: util.Map[String, Array[Byte]] = namespaceAwareCoordinatorStreamStore.all
val configMap: util.Map[String, String] = new util.HashMap[String, String]
for ((key: String, valueAsBytes: Array[Byte]) <- configFromCoordinatorStream.asScala) {
val valueSerde: CoordinatorStreamValueSerde = new CoordinatorStreamValueSerde(SetConfig.TYPE)
val valueAsString: String = valueSerde.fromBytes(valueAsBytes)
configMap.put(key, valueAsString)
if (valueAsBytes == null) {
warn("Value for key: %s in config is null. Ignoring it." format key)
} 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)
} else {
configMap.put(key, valueAsString)
}
}
}
new MapConfig(configMap)
}
Expand Down