From 07918b32977344bedd1200dae2f66370bbf271c5 Mon Sep 17 00:00:00 2001 From: Shanthoosh Venkataraman Date: Fri, 26 Apr 2019 15:25:26 -0700 Subject: [PATCH] SAMZA-2176: Ignore the configuration keys with serialized null values from the coordinator stream. --- .../samza/util/CoordinatorStreamUtil.scala | 17 +++++++++++++---- 1 file changed, 13 insertions(+), 4 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 2fca36c6f3..485a55c004 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 @@ -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} @@ -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. @@ -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) }