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
12 changes: 12 additions & 0 deletions samza-core/src/main/java/org/apache/samza/config/JobConfig.java
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,9 @@ public class JobConfig extends MapConfig {
public static final String CONTAINER_METADATA_FILENAME_FORMAT = "%s.metadata"; // Filename: <containerID>.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);
}
Expand Down Expand Up @@ -341,4 +344,13 @@ public static Optional<File> 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);
}

}
Original file line number Diff line number Diff line change
@@ -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);
}
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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}
Expand All @@ -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 {
/**
Expand All @@ -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);

}

/**
Expand Down
Original file line number Diff line number Diff line change
@@ -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 {
Comment thread
MabelYC marked this conversation as resolved.
@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<String, String> map = config.subset(String.format(SystemConfig.SYSTEM_ID_PREFIX, jobConfig.getCoordinatorSystemName()), false);
Map<String, String> 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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
}
Original file line number Diff line number Diff line change
@@ -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<String, String> 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<String, String> 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<String, String> 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));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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 {

Expand All @@ -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)

Comment thread
xinyuiscool marked this conversation as resolved.
Assert.assertEquals(configMap.get("systems.samzatest.test"), "test")
Assert.assertEquals(configMap.get("test.only"), null)
}

@Test
def testReadConfigFromCoordinatorStream {
val keyForNonBlankVal = "app.id"
Expand Down