From 5448c1fe6d671834240c79703e20e0d4cd62fd31 Mon Sep 17 00:00:00 2001 From: Aditya <60684641+0x0elliot@users.noreply.github.com> Date: Fri, 5 Jul 2024 18:12:51 +0530 Subject: [PATCH] fix[k8s]: worker replication works at scale --- functions/onprem/orborus/orborus.go | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index 02a57740..ffc306e6 100755 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -831,6 +831,9 @@ func fixk8sRoles() { } } + +func int32Ptr(i int32) *int32 { return &i } + func deployK8sWorker(image string, identifier string, env []string) error { env = append(env, fmt.Sprintf("IS_KUBERNETES=true")) env = append(env, fmt.Sprintf("KUBERNETES_NAMESPACE=%s", os.Getenv("KUBERNETES_NAMESPACE"))) @@ -1015,11 +1018,25 @@ func deployK8sWorker(image string, identifier string, env []string) error { // return err // } + replicaNumberStr := os.Getenv("SHUFFLE_SCALE_REPLICAS") + replicaNumber := 1 + if len(replicaNumberStr) > 0 { + tmpInt, err := strconv.Atoi(replicaNumberStr) + if err != nil { + log.Printf("[ERROR] %s is not a valid number for replication", replicaNumberStr) + } else { + replicaNumber = tmpInt + } + } + + replicaNumberInt32 := int32(replicaNumber) + deployment := &appsv1.Deployment{ ObjectMeta: metav1.ObjectMeta{ Name: identifier, }, Spec: appsv1.DeploymentSpec{ + Replicas: int32Ptr(replicaNumberInt32), Selector: &metav1.LabelSelector{ MatchLabels: containerLabels, },