diff --git a/cdap-app-fabric/src/main/java/io/cdap/cdap/internal/app/worker/TaskWorkerServiceLauncher.java b/cdap-app-fabric/src/main/java/io/cdap/cdap/internal/app/worker/TaskWorkerServiceLauncher.java index 8f43c9697366..de895d352cda 100644 --- a/cdap-app-fabric/src/main/java/io/cdap/cdap/internal/app/worker/TaskWorkerServiceLauncher.java +++ b/cdap-app-fabric/src/main/java/io/cdap/cdap/internal/app/worker/TaskWorkerServiceLauncher.java @@ -203,8 +203,29 @@ public void run() { Map configMap = new HashMap<>(); configMap.put(ProgramOptionConstants.RUNTIME_NAMESPACE, NamespaceId.SYSTEM.getNamespace()); + if (cConf.get(Constants.TaskWorker.CONTAINER_CPU_MULTIPLIER) != null) { + configMap.put(Constants.Kube.CPU_MULTIPLIER, + cConf.get(Constants.TaskWorker.CONTAINER_CPU_MULTIPLIER)); + } + if (cConf.get(Constants.TaskWorker.CONTAINER_MEMORY_MULTIPLIER) != null) { + configMap.put(Constants.Kube.MEMORY_MULTIPLIER, + cConf.get(Constants.TaskWorker.CONTAINER_MEMORY_MULTIPLIER)); + } twillPreparer.withConfiguration(Collections.unmodifiableMap(configMap)); + Map artifactLocalizerConfig = new HashMap<>(); + if (cConf.get(Constants.ArtifactLocalizer.CONTAINER_CPU_MULTIPLIER) != null) { + artifactLocalizerConfig.put(Constants.Kube.CPU_MULTIPLIER, + cConf.get(Constants.ArtifactLocalizer.CONTAINER_CPU_MULTIPLIER)); + } + if (cConf.get(Constants.ArtifactLocalizer.CONTAINER_MEMORY_MULTIPLIER) != null) { + artifactLocalizerConfig.put(Constants.Kube.MEMORY_MULTIPLIER, + cConf.get(Constants.ArtifactLocalizer.CONTAINER_MEMORY_MULTIPLIER)); + } + if (!artifactLocalizerConfig.isEmpty()) { + twillPreparer.withConfiguration("ArtifactLocalizerTwillRunnable", artifactLocalizerConfig); + } + if (Feature.NAMESPACED_SERVICE_ACCOUNTS.isEnabled(featureFlagsProvider)) { String localhost = InetAddress.getLoopbackAddress().getHostName(); twillPreparer = twillPreparer.withEnv(TaskWorkerTwillRunnable.class.getSimpleName(), diff --git a/cdap-common/src/main/java/io/cdap/cdap/common/conf/Constants.java b/cdap-common/src/main/java/io/cdap/cdap/common/conf/Constants.java index ccb27ae8a6f8..e8e28549a43f 100644 --- a/cdap-common/src/main/java/io/cdap/cdap/common/conf/Constants.java +++ b/cdap-common/src/main/java/io/cdap/cdap/common/conf/Constants.java @@ -490,6 +490,14 @@ public static final class Environment { */ public static final String PROGRAM_SUBMISSION_MASTER_ENV_ENABLED = "program.submission.master.environment.enabled"; } + /** + * Kubernetes constants. + * These same as defined in KubeTwillPreparer + */ + public static final class Kube { + public static final String CPU_MULTIPLIER = "master.environment.k8s.container.cpu.multiplier"; + public static final String MEMORY_MULTIPLIER = "master.environment.k8s.container.memory.multiplier"; + } /** * Task worker. @@ -610,6 +618,8 @@ public static final class ArtifactLocalizer { public static final String CONTAINER_MEMORY_MB = "artifact.localizer.container.memory.mb"; public static final String CONTAINER_CORES = "artifact.localizer.container.num.cores"; public static final String CONTAINER_JVM_OPTS = "artifact.localizer.container.jvm.opts"; + public static final String CONTAINER_CPU_MULTIPLIER = "artifact.localizer.container.cpu.multiplier"; + public static final String CONTAINER_MEMORY_MULTIPLIER = "artifact.localizer.container.memory.multiplier"; /** * Artifact localizer http handler configuration. diff --git a/cdap-kubernetes/src/main/java/io/cdap/cdap/k8s/runtime/KubeTwillPreparer.java b/cdap-kubernetes/src/main/java/io/cdap/cdap/k8s/runtime/KubeTwillPreparer.java index 2025cd7076ae..df6db9c78a4d 100644 --- a/cdap-kubernetes/src/main/java/io/cdap/cdap/k8s/runtime/KubeTwillPreparer.java +++ b/cdap-kubernetes/src/main/java/io/cdap/cdap/k8s/runtime/KubeTwillPreparer.java @@ -215,6 +215,8 @@ class KubeTwillPreparer implements DependentTwillPreparer, StatefulTwillPreparer private StringBuilder globalJvmOptions; private final V1EmptyDirVolumeSource workDirVolumeSource; private boolean shouldLocalizeConfigurationAsConfigmap; + private String systemCpuMultiplier; + private String systemMemoryMultiplier; KubeTwillPreparer(MasterEnvironmentContext masterEnvContext, ApiClient apiClient, String kubeNamespace, @@ -436,6 +438,12 @@ public TwillPreparer withConfiguration(Map config) { if (config.containsKey(MasterOptionConstants.RUNTIME_NAMESPACE)) { cdapRuntimeNamespace = config.get(MasterOptionConstants.RUNTIME_NAMESPACE); } + if (config.containsKey(CPU_MULTIPLIER)) { + systemCpuMultiplier = config.get(CPU_MULTIPLIER); + } + if (config.containsKey(MEMORY_MULTIPLIER)) { + systemMemoryMultiplier = config.get(MEMORY_MULTIPLIER); + } for (String runnableName : runnables) { withEnv(runnableName, config); } @@ -1201,7 +1209,7 @@ private V1PodSpec createPodSpec(Location runtimeConfigLocation, RuntimeSpecification mainRuntimeSpec = getMainRuntimeSpecification(runtimeSpecs); String runnableName = mainRuntimeSpec.getName(); final V1ResourceRequirements initContainerResourceRequirements = - createResourceRequirements(mainRuntimeSpec.getResourceSpecification()); + createResourceRequirements(mainRuntimeSpec.getName(), mainRuntimeSpec.getResourceSpecification()); // Setup the container environment. Inherit everything from the current pod except workload identity env vars. Map initContainerEnvirons = podInfo.getContainerEnvironments().stream() @@ -1330,7 +1338,7 @@ private List createContainers(Map run containers.add(createContainer(mainRuntimeSpec.getName(), podInfo.getContainerImage(), podInfo.getImagePullPolicy(), workDir, - createResourceRequirements(mainRuntimeSpec.getResourceSpecification()), + createResourceRequirements(mainRuntimeSpec.getName(), mainRuntimeSpec.getResourceSpecification()), mounts, environs, KubeTwillLauncher.class, Stream.concat(Stream.of(mainRuntimeSpec.getName()), args.stream()) .toArray(String[]::new))); @@ -1356,7 +1364,7 @@ private List createContainers(Map run mounts = addSecreteVolMountIfNeeded(spec, volumeMounts); containers.add( createContainer(name, podInfo.getContainerImage(), podInfo.getImagePullPolicy(), workDir, - createResourceRequirements(spec.getResourceSpecification()), + createResourceRequirements(name, spec.getResourceSpecification()), mounts, envs, KubeTwillLauncher.class, Stream.concat(Stream.of(name), args.stream()).toArray(String[]::new))); } @@ -1452,11 +1460,23 @@ private V1Container createContainer(String name, String containerImage, String i * the namespace has a resource quota, the objects must also specify resource limits. */ @VisibleForTesting - V1ResourceRequirements createResourceRequirements(ResourceSpecification resourceSpec) { + V1ResourceRequirements createResourceRequirements(String runnableName, ResourceSpecification resourceSpec) { Map cConf = masterEnvContext.getConfigurations(); - float cpuMultiplier = Float.parseFloat(cConf.getOrDefault(CPU_MULTIPLIER, DEFAULT_MULTIPLIER)); - float memoryMultiplier = Float.parseFloat( - cConf.getOrDefault(MEMORY_MULTIPLIER, DEFAULT_MULTIPLIER)); + + String runnableCpu = null; + String runnableMem = null; + if (runnableName != null && runnableConfigs.containsKey(runnableName)) { + Map configMap = runnableConfigs.get(runnableName); + runnableCpu = configMap.get(CPU_MULTIPLIER); + runnableMem = configMap.get(MEMORY_MULTIPLIER); + } + + float cpuMultiplier = Float.parseFloat(runnableCpu != null ? runnableCpu : + (systemCpuMultiplier != null ? systemCpuMultiplier : + cConf.getOrDefault(CPU_MULTIPLIER, DEFAULT_MULTIPLIER))); + float memoryMultiplier = Float.parseFloat(runnableMem != null ? runnableMem : + (systemMemoryMultiplier != null ? systemMemoryMultiplier : + cConf.getOrDefault(MEMORY_MULTIPLIER, DEFAULT_MULTIPLIER))); V1ResourceRequirementsBuilder requirementsBuilder = new V1ResourceRequirementsBuilder(); diff --git a/cdap-kubernetes/src/test/java/io/cdap/cdap/k8s/runtime/KubeTwillPreparerTest.java b/cdap-kubernetes/src/test/java/io/cdap/cdap/k8s/runtime/KubeTwillPreparerTest.java index e8ccacc6d8fa..6e224e0661b5 100644 --- a/cdap-kubernetes/src/test/java/io/cdap/cdap/k8s/runtime/KubeTwillPreparerTest.java +++ b/cdap-kubernetes/src/test/java/io/cdap/cdap/k8s/runtime/KubeTwillPreparerTest.java @@ -168,7 +168,7 @@ public void testCreateDefaultResourceSpecification() throws Exception { Map config = new HashMap<>(); config.put(MasterOptionConstants.RUNTIME_NAMESPACE, "system"); preparer.withConfiguration(config); - V1ResourceRequirements gotResourceRequirements = preparer.createResourceRequirements(resourceSpecification); + V1ResourceRequirements gotResourceRequirements = preparer.createResourceRequirements(null, resourceSpecification); Assert.assertEquals("1", gotResourceRequirements.getRequests().get("cpu").toSuffixedString()); Assert.assertEquals("100Mi", gotResourceRequirements.getRequests().get("memory").toSuffixedString()); } @@ -183,7 +183,7 @@ public void testCreateDefaultSystemResourceSpecification() throws Exception { preparer.withConfiguration(config); ResourceSpecification resourceSpecification = new DefaultResourceSpecification(1, 100, 1, 1, 1); - V1ResourceRequirements gotResourceRequirements = preparer.createResourceRequirements(resourceSpecification); + V1ResourceRequirements gotResourceRequirements = preparer.createResourceRequirements(null, resourceSpecification); Assert.assertEquals("1", gotResourceRequirements.getRequests().get("cpu").toSuffixedString()); Assert.assertEquals("100Mi", gotResourceRequirements.getRequests().get("memory").toSuffixedString()); } @@ -225,7 +225,7 @@ public void testCreateResourceSpecificationWithCustomResourceMultipliers() throw config.put(MasterOptionConstants.RUNTIME_NAMESPACE, "system"); preparer.withConfiguration(config); ResourceSpecification resourceSpecification = new DefaultResourceSpecification(1, 100, 1, 1, 1); - V1ResourceRequirements gotResourceRequirements = preparer.createResourceRequirements(resourceSpecification); + V1ResourceRequirements gotResourceRequirements = preparer.createResourceRequirements(null, resourceSpecification); Assert.assertEquals("500m", gotResourceRequirements.getRequests().get("cpu").toSuffixedString()); Assert.assertEquals("25Mi", gotResourceRequirements.getRequests().get("memory").toSuffixedString()); } @@ -241,7 +241,7 @@ public void testCreateDefaultUserResourceSpecification() throws Exception { preparer.withConfiguration(config); ResourceSpecification resourceSpecification = new DefaultResourceSpecification(1, 100, 1, 1, 1); - V1ResourceRequirements gotResourceRequirements = preparer.createResourceRequirements(resourceSpecification); + V1ResourceRequirements gotResourceRequirements = preparer.createResourceRequirements(null, resourceSpecification); Assert.assertEquals("500m", gotResourceRequirements.getRequests().get("cpu").toSuffixedString()); Assert.assertEquals("50Mi", gotResourceRequirements.getRequests().get("memory").toSuffixedString()); Assert.assertEquals("1", gotResourceRequirements.getLimits().get("cpu").toSuffixedString()); @@ -261,7 +261,7 @@ public void testCreateUserResourceSpecificationWithCustomResourceMultipliers() t preparer.withConfiguration(config); ResourceSpecification resourceSpecification = new DefaultResourceSpecification(1, 100, 1, 1, 1); - V1ResourceRequirements gotResourceRequirements = preparer.createResourceRequirements(resourceSpecification); + V1ResourceRequirements gotResourceRequirements = preparer.createResourceRequirements(null, resourceSpecification); Assert.assertEquals("300m", gotResourceRequirements.getRequests().get("cpu").toSuffixedString()); Assert.assertEquals("70Mi", gotResourceRequirements.getRequests().get("memory").toSuffixedString()); Assert.assertEquals("1", gotResourceRequirements.getLimits().get("cpu").toSuffixedString()); @@ -280,7 +280,7 @@ public void testCreateUserResourceSpecificationInvalidProgramCpuMultiplier() thr preparer.withConfiguration(config); ResourceSpecification resourceSpecification = new DefaultResourceSpecification(1, 100, 1, 1, 1); - preparer.createResourceRequirements(resourceSpecification); + preparer.createResourceRequirements(null, resourceSpecification); } @Test(expected = IllegalArgumentException.class) @@ -298,7 +298,7 @@ public void testCreateUserResourceSpecificationInvalidProgramMemoryMultiplier() ResourceSpecification resourceSpecification = new DefaultResourceSpecification( 1, 100, 1, 1, 1); - preparer.createResourceRequirements(resourceSpecification); + preparer.createResourceRequirements(null, resourceSpecification); } @Test diff --git a/cdap-master/src/main/java/io/cdap/cdap/master/environment/k8s/PreviewServiceMain.java b/cdap-master/src/main/java/io/cdap/cdap/master/environment/k8s/PreviewServiceMain.java index c0c3c4628641..577c3b1aeced 100644 --- a/cdap-master/src/main/java/io/cdap/cdap/master/environment/k8s/PreviewServiceMain.java +++ b/cdap-master/src/main/java/io/cdap/cdap/master/environment/k8s/PreviewServiceMain.java @@ -65,10 +65,6 @@ */ public class PreviewServiceMain extends AbstractServiceMain { - // Following constants are same as defined in AbstractKubeTwillPreparer - private static final String KUBE_CPU_MULTIPLIER = "master.environment.k8s.container.cpu.multiplier"; - private static final String KUBE_MEMORY_MULTIPLIER = "master.environment.k8s.container.memory.multiplier"; - /** * Main entry point */ @@ -79,8 +75,8 @@ public static void main(String[] args) throws Exception { @Override protected CConfiguration updateCConf(CConfiguration cConf) { Map keyMap = ImmutableMap.of( - Constants.Preview.CONTAINER_CPU_MULTIPLIER, KUBE_CPU_MULTIPLIER, - Constants.Preview.CONTAINER_MEMORY_MULTIPLIER, KUBE_MEMORY_MULTIPLIER, + Constants.Preview.CONTAINER_CPU_MULTIPLIER, Constants.Kube.CPU_MULTIPLIER, + Constants.Preview.CONTAINER_MEMORY_MULTIPLIER, Constants.Kube.MEMORY_MULTIPLIER, Constants.Preview.CONTAINER_HEAP_RESERVED_RATIO, Configs.Keys.HEAP_RESERVED_MIN_RATIO );