diff --git a/samza-core/src/main/java/org/apache/samza/config/JobConfig.java b/samza-core/src/main/java/org/apache/samza/config/JobConfig.java index 9bdcfd5cfa..efff5724b6 100644 --- a/samza-core/src/main/java/org/apache/samza/config/JobConfig.java +++ b/samza-core/src/main/java/org/apache/samza/config/JobConfig.java @@ -125,6 +125,9 @@ public class JobConfig extends MapConfig { public static final String CONTAINER_METADATA_FILENAME_FORMAT = "%s.metadata"; // Filename: .metadata public static final String CONTAINER_METADATA_DIRECTORY_SYS_PROPERTY = "samza.log.dir"; + public static final String COORDINATOR_STREAM_FACTORY = "job.coordinatorstream.config.factory"; + public static final String DEFAULT_COORDINATOR_STREAM_CONFIG_FACTORY = "org.apache.samza.util.DefaultCoordinatorStreamConfigFactory"; + public JobConfig(Config config) { super(config); } @@ -341,4 +344,13 @@ public static Optional getMetadataFile(String execEnvContainerId) { new File(dir, String.format(CONTAINER_METADATA_FILENAME_FORMAT, execEnvContainerId))); } } + + /** + * Get coordinatorStreamFactory according to the configs + * @return the name of coordinatorStreamFactory + */ + public String getCoordinatorStreamFactory() { + return get(COORDINATOR_STREAM_FACTORY, DEFAULT_COORDINATOR_STREAM_CONFIG_FACTORY); + } + } \ No newline at end of file diff --git a/samza-core/src/main/scala/org/apache/samza/util/CoordinatorStreamConfigFactory.java b/samza-core/src/main/scala/org/apache/samza/util/CoordinatorStreamConfigFactory.java new file mode 100644 index 0000000000..6f930f3867 --- /dev/null +++ b/samza-core/src/main/scala/org/apache/samza/util/CoordinatorStreamConfigFactory.java @@ -0,0 +1,38 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.samza.util; + +import org.apache.samza.config.Config; + + +/** + * A CoordinatorStreamConfigFactory receives the job's config and create specific configs that are needed to + * create coordinator streams. + */ +public interface CoordinatorStreamConfigFactory { + + /** + * Returns a Config what is needed to create coordinator streams. + * + * @param config known configs for job + * @return basic configs needed to create coordinator streams + */ + Config buildCoordinatorStreamConfig(Config config); +} 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 e38359db9f..f6636e7663 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 @@ -1,5 +1,4 @@ /* - * * Licensed to the Apache Software Foundation (ASF) under one * or more contributor license agreements. See the NOTICE file * distributed with this work for additional information @@ -23,7 +22,6 @@ 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.coordinator.metadatastore.{CoordinatorStreamStore, NamespaceAwareCoordinatorStreamStore} @@ -33,7 +31,6 @@ import org.apache.samza.system.{StreamSpec, SystemAdmin, SystemFactory, SystemSt import org.apache.samza.util.ScalaJavaUtil.JavaOptionals import scala.collection.JavaConverters._ -import scala.collection.immutable.Map object CoordinatorStreamUtil extends Logging { /** @@ -42,15 +39,11 @@ object CoordinatorStreamUtil extends Logging { */ def buildCoordinatorStreamConfig(config: Config) = { val jobConfig = new JobConfig(config) - val (jobName, jobId) = getJobNameAndId(jobConfig) - // Build a map with just the system config and job.name/job.id. This is what's required to start the JobCoordinator. - val map = config.subset(SystemConfig.SYSTEM_ID_PREFIX format jobConfig.getCoordinatorSystemName, false).asScala ++ - Map[String, String]( - JobConfig.JOB_NAME -> jobName, - JobConfig.JOB_ID -> jobId, - JobConfig.JOB_COORDINATOR_SYSTEM -> jobConfig.getCoordinatorSystemName, - JobConfig.MONITOR_PARTITION_CHANGE_FREQUENCY_MS -> String.valueOf(jobConfig.getMonitorPartitionChangeFrequency)) - new MapConfig(map.asJava) + val buildConfigFactory = jobConfig.getCoordinatorStreamFactory(); + val coordinatorSystemConfig = Class.forName(buildConfigFactory).newInstance().asInstanceOf[CoordinatorStreamConfigFactory].buildCoordinatorStreamConfig(config) + + new MapConfig(coordinatorSystemConfig); + } /** diff --git a/samza-core/src/main/scala/org/apache/samza/util/DefaultCoordinatorStreamConfigFactory.java b/samza-core/src/main/scala/org/apache/samza/util/DefaultCoordinatorStreamConfigFactory.java new file mode 100644 index 0000000000..f3bf10faba --- /dev/null +++ b/samza-core/src/main/scala/org/apache/samza/util/DefaultCoordinatorStreamConfigFactory.java @@ -0,0 +1,50 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.samza.util; + +import java.util.HashMap; +import java.util.Map; +import org.apache.samza.config.Config; +import org.apache.samza.config.ConfigException; +import org.apache.samza.config.JobConfig; +import org.apache.samza.config.MapConfig; +import org.apache.samza.config.SystemConfig; + + +public class DefaultCoordinatorStreamConfigFactory implements CoordinatorStreamConfigFactory { + @Override + public Config buildCoordinatorStreamConfig(Config config) { + JobConfig jobConfig = new JobConfig(config); + String jobName = jobConfig.getName().orElseThrow(() -> new ConfigException("Missing required config: job.name")); + + String jobId = jobConfig.getJobId(); + + // Build a map with just the system config and job.name/job.id. This is what's required to start the JobCoordinator. + Map map = config.subset(String.format(SystemConfig.SYSTEM_ID_PREFIX, jobConfig.getCoordinatorSystemName()), false); + Map addConfig = new HashMap<>(); + addConfig.put(JobConfig.JOB_NAME, jobName); + addConfig.put(JobConfig.JOB_ID, jobId); + addConfig.put(JobConfig.JOB_COORDINATOR_SYSTEM, jobConfig.getCoordinatorSystemName()); + addConfig.put(JobConfig.MONITOR_PARTITION_CHANGE_FREQUENCY_MS, String.valueOf(jobConfig.getMonitorPartitionChangeFrequency())); + + addConfig.putAll(map); + return new MapConfig(addConfig); + } +} diff --git a/samza-core/src/test/java/org/apache/samza/config/TestJobConfig.java b/samza-core/src/test/java/org/apache/samza/config/TestJobConfig.java index cf14ae3851..dab0f77ef6 100644 --- a/samza-core/src/test/java/org/apache/samza/config/TestJobConfig.java +++ b/samza-core/src/test/java/org/apache/samza/config/TestJobConfig.java @@ -576,4 +576,13 @@ public void testGetMetadataFile() { assertEquals(Optional.empty(), JobConfig.getMetadataFile(null)); } + + @Test + public void testGetCoordinatorStreamFactory() { + JobConfig jobConfig = new JobConfig(new MapConfig(ImmutableMap.of("test", ""))); + assertEquals(jobConfig.getCoordinatorStreamFactory(), JobConfig.DEFAULT_COORDINATOR_STREAM_CONFIG_FACTORY); + + jobConfig = new JobConfig(new MapConfig(ImmutableMap.of(JobConfig.COORDINATOR_STREAM_FACTORY, "specific_coordinator_stream"))); + assertEquals(jobConfig.getCoordinatorStreamFactory(), "specific_coordinator_stream"); + } } \ No newline at end of file diff --git a/samza-core/src/test/java/org/apache/samza/util/TestDefaultCoordinatorStreamConfigFactory.java b/samza-core/src/test/java/org/apache/samza/util/TestDefaultCoordinatorStreamConfigFactory.java new file mode 100644 index 0000000000..209dd05a3e --- /dev/null +++ b/samza-core/src/test/java/org/apache/samza/util/TestDefaultCoordinatorStreamConfigFactory.java @@ -0,0 +1,69 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + * + */ + +package org.apache.samza.util; + +import java.util.HashMap; +import java.util.Map; +import org.apache.samza.config.Config; +import org.apache.samza.config.ConfigException; +import org.apache.samza.config.JobConfig; +import org.apache.samza.config.MapConfig; +import org.junit.Test; + +import static org.junit.Assert.*; + + +public class TestDefaultCoordinatorStreamConfigFactory { + + DefaultCoordinatorStreamConfigFactory factory = new DefaultCoordinatorStreamConfigFactory(); + + @Test + public void testBuildCoordinatorStreamConfigWithJobName() { + Map mapConfig = new HashMap<>(); + mapConfig.put("job.name", "testName"); + mapConfig.put("job.id", "testId"); + mapConfig.put("job.coordinator.system", "testSamza"); + mapConfig.put("test.only", "nothing"); + mapConfig.put("systems.testSamza.test", "test"); + + Config config = factory.buildCoordinatorStreamConfig(new MapConfig(mapConfig)); + + Map expectedMap = new HashMap<>(); + expectedMap.put("job.name", "testName"); + expectedMap.put("job.id", "testId"); + expectedMap.put("systems.testSamza.test", "test"); + expectedMap.put(JobConfig.JOB_COORDINATOR_SYSTEM, "testSamza"); + expectedMap.put(JobConfig.MONITOR_PARTITION_CHANGE_FREQUENCY_MS, "300000"); + + assertEquals(config, new MapConfig(expectedMap)); + } + + @Test(expected = ConfigException.class) + public void testBuildCoordinatorStreamConfigWithoutJobName() { + Map mapConfig = new HashMap<>(); + mapConfig.put("job.id", "testId"); + mapConfig.put("job.coordinator.system", "testSamza"); + mapConfig.put("test.only", "nothing"); + mapConfig.put("systems.testSamza.test", "test"); + + factory.buildCoordinatorStreamConfig(new MapConfig(mapConfig)); + } +} 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 eb035960d9..dac1fe08a8 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 @@ -27,6 +27,7 @@ import org.apache.samza.system.{StreamSpec, SystemAdmin, SystemStream} import org.junit.{Assert, Test} import org.mockito.Matchers.any import org.mockito.Mockito +import org.apache.samza.config.MapConfig class TestCoordinatorStreamUtil { @@ -40,6 +41,21 @@ class TestCoordinatorStreamUtil { Mockito.verify(systemAdmin).createStream(any(classOf[StreamSpec])) } + @Test + def testBuildCoordinatorStreamConfig: Unit = { + val addConfig = new util.HashMap[String, String] + addConfig.put("job.name", "test-job-name") + addConfig.put("job.id", "i001") + addConfig.put("job.coordinator.system", "samzatest") + addConfig.put("systems.samzatest.test","test") + addConfig.put("test.only","nothing") + val config = new MapConfig(addConfig) + val configMap = CoordinatorStreamUtil.buildCoordinatorStreamConfig(config) + + Assert.assertEquals(configMap.get("systems.samzatest.test"), "test") + Assert.assertEquals(configMap.get("test.only"), null) + } + @Test def testReadConfigFromCoordinatorStream { val keyForNonBlankVal = "app.id"