made orborus use time window functions and change the scaling per worker ran requests

This commit is contained in:
yashsinghcodes
2024-11-30 02:50:02 +05:30
parent a7af26d28e
commit 2e715efbfd
+59 -78
View File
@@ -38,7 +38,6 @@ import (
"github.com/docker/docker/api/types" "github.com/docker/docker/api/types"
"github.com/docker/docker/api/types/container" "github.com/docker/docker/api/types/container"
"github.com/docker/docker/api/types/filters"
"github.com/docker/docker/api/types/image" "github.com/docker/docker/api/types/image"
"github.com/docker/docker/api/types/mount" "github.com/docker/docker/api/types/mount"
"github.com/docker/docker/api/types/network" "github.com/docker/docker/api/types/network"
@@ -105,7 +104,7 @@ var swarmNetworkName = os.Getenv("SHUFFLE_SWARM_NETWORK_NAME")
var orborusLabel = os.Getenv("SHUFFLE_ORBORUS_LABEL") var orborusLabel = os.Getenv("SHUFFLE_ORBORUS_LABEL")
var memcached = os.Getenv("SHUFFLE_MEMCACHED") var memcached = os.Getenv("SHUFFLE_MEMCACHED")
var queuePerMinute = os.Getenv("SHUFFLE_EXECUTION_PER_MINIUTE") var queuePerMinute = os.Getenv("SHUFFLE_EXECUTION_PER_MINIUTE")
var queuePerMinuteInt = queuePerMinute var queuePerMinuteInt int
// For it to download from Sigma? // For it to download from Sigma?
var apiKey = os.Getenv("AUTH_FOR_ORBORUS") var apiKey = os.Getenv("AUTH_FOR_ORBORUS")
@@ -122,6 +121,7 @@ var containerId string
var executionCount = 0 var executionCount = 0
var imagedownloadTimeout = time.Second * 300 var imagedownloadTimeout = time.Second * 300
var window = shuffle.NewTimeWindow(1 * time.Minute)
func init() { func init() {
var err error var err error
@@ -330,8 +330,9 @@ func deployServiceWorkers(image string) {
log.Printf("[DEBUG] Skipping deployment of workers as services as swarmConfig is not set to run or swarm. Value: %#v", swarmConfig) log.Printf("[DEBUG] Skipping deployment of workers as services as swarmConfig is not set to run or swarm. Value: %#v", swarmConfig)
return return
} }
ctx := context.Background()
isMemcachedRunning, err := checkMemcached(dockercli) isMemcachedRunning, err := checkMemcached(ctx, dockercli)
if err != nil { if err != nil {
log.Printf("[ERROR] Failed checking memcached: %s", err) log.Printf("[ERROR] Failed checking memcached: %s", err)
} }
@@ -343,7 +344,6 @@ func deployServiceWorkers(image string) {
os.Setenv("SHUFFLE_MEMCACHED", fmt.Sprintf("%s:11211", ip)) os.Setenv("SHUFFLE_MEMCACHED", fmt.Sprintf("%s:11211", ip))
ctx := context.Background()
// Looks for and cleans up all existing items in swarm we can't re-use (Shuffle only) // Looks for and cleans up all existing items in swarm we can't re-use (Shuffle only)
// frikky@debian:~/git/shuffle/functions/onprem/worker$ docker service create --replicas 5 --name shuffle-workers --env SHUFFLE_SWARM_CONFIG=run --publish published=33333,target=33333 ghcr.io/shuffle/shuffle-worker:nightly // frikky@debian:~/git/shuffle/functions/onprem/worker$ docker service create --replicas 5 --name shuffle-workers --env SHUFFLE_SWARM_CONFIG=run --publish published=33333,target=33333 ghcr.io/shuffle/shuffle-worker:nightly
@@ -3797,6 +3797,7 @@ func sendWorkerRequest(workflowExecution shuffle.ExecutionRequest, image string,
log.Printf("[ERROR] Failed reading body in worker request body to worker on %s: %s", streamUrl, err) log.Printf("[ERROR] Failed reading body in worker request body to worker on %s: %s", streamUrl, err)
return err return err
} }
window.AddEvent(time.Now())
if newresp.StatusCode != 200 { if newresp.StatusCode != 200 {
log.Printf("[WARNING] POTENTIAL error running worker request (2) - status code is %d for %s, not 200. Body: %s", newresp.StatusCode, streamUrl, string(body)) log.Printf("[WARNING] POTENTIAL error running worker request (2) - status code is %d for %s, not 200. Body: %s", newresp.StatusCode, streamUrl, string(body))
@@ -3825,85 +3826,43 @@ func AutoScale(ctx context.Context) {
return return
} }
if queuePerMinute != "" { ticker := time.NewTicker(1 * time.Second)
coolDownPeriod := 10 * time.Second
queuePerMinuteInt = 20
if os.Getenv("SHUFFLE_QUEUE_PER_MINUTE") != "" {
var err error var err error
queuePerMinuteInt, err = strconv.Atoi(queuePerMinute) queuePerMinuteInt, err = strconv.Atoi(os.Getenv("SHUFFLE_QUEUE_PER_MINUTE"))
if err != nil { if err != nil {
queuePerMinuteInt = 60 log.Printf("[WARNING] Cannot convert %s to int. Using default value for it: %d", queuePerMinute, queuePerMinuteInt)
log.Printf("[WARNING] Cannot convert %s to int please pass an interger value. Using default value for it %d", queuePerMinute, queuePerMinuteInt)
} }
} }
replicas := uint64(6) lastScaleTime := time.Now()
scaleReplicas := os.Getenv("SHUFFLE_SCALE_REPLICAS") currentWorkers := currentWokerCount(ctx, dockercli)
if len(scaleReplicas) > 0 {
tmpInt, err := strconv.Atoi(scaleReplicas)
if err != nil {
log.Printf("[ERROR] %s is not a valid number for replication", scaleReplicas)
} else {
replicas = uint64(tmpInt)
}
log.Printf("[DEBUG] SHUFFLE_SCALE_REPLICAS set to value %#v. Trying to overwrite default (%d/node)", scaleReplicas, replicas)
}
config := shuffle.ScalingConfig{
QueuePerMinute: queuePerMinuteInt,
ScalingInterval: 12 * time.Second, // TO ADD SOME BUFFER
MaxScaleUpStep: 1,
CooldownPeriod: time.Duration(10 * time.Second),
MaxReplicas: int(replicas),
}
lastScaleTime := time.Now().Add(-config.CooldownPeriod)
var lastQueueLength int
var lastCheckTime time.Time
for { for {
select { select {
case <-ctx.Done(): case <-ctx.Done():
return return
case <-time.After(config.ScalingInterval): case <-ticker.C:
if time.Since(lastScaleTime) < config.CooldownPeriod { if time.Since(lastScaleTime) < (coolDownPeriod) {
continue continue
} }
currentRequestCount := window.CountEvents(time.Now())
data, err := collectMetrics(ctx, dockercli) requiredReplicas := 0
if err != nil { if currentRequestCount >= queuePerMinuteInt*currentWorkers {
continue // FIXME: Hardcoded Max Replicas should be 6
requiredReplicas = int(math.Min(float64(6), float64(currentRequestCount/queuePerMinuteInt)+1))
} }
currentTime := time.Now() if requiredReplicas > 0 {
currentQueueLength := data err := scaleService(ctx, dockercli, uint64(requiredReplicas))
if err != nil {
if !lastCheckTime.IsZero() { log.Printf("[ERROR] Failed to scale service: %s", err)
timeDiff := currentTime.Sub(lastCheckTime).Seconds() } else {
if timeDiff >= 10 { lastScaleTime = time.Now()
rocq := int(math.Ceil(float64(currentQueueLength - lastQueueLength/int(timeDiff)))) currentWorkers = currentWokerCount(ctx, dockercli)
//desiredReplicas, currentReplicas := numberOfReplicas(ctx, data, config)
requiredReplicas := 0
if rocq >= config.QueuePerMinute {
requiredReplicas = currentQueueLength + config.MaxScaleUpStep
if requiredReplicas > config.MaxReplicas {
requiredReplicas = config.MaxReplicas
}
}
if requiredReplicas != 0 {
err := scaleService(ctx, dockercli, uint64(requiredReplicas))
if err != nil {
log.Printf("[ERROR] Failed to scale the service: %s", err)
} else {
lastScaleTime = currentTime
}
}
} }
} }
lastQueueLength = currentQueueLength
lastCheckTime = currentTime
} }
} }
} }
@@ -3933,16 +3892,29 @@ func scaleService(ctx context.Context, client *dockerclient.Client, replicas uin
return nil return nil
} }
func queueScaleFactor(numQueue int, config shuffle.ScalingConfig) float64 { func currentWokerCount(ctx context.Context, client *dockerclient.Client) int {
if numQueue > config.QueuePerMinute { service, _, err := client.ServiceInspectWithRaw(ctx, "shuffle-workers", types.ServiceInspectOptions{})
queuePressure := float64(numQueue) / float64(config.QueuePerMinute) if err != nil {
return 0
}
if service.Spec.Mode.Replicated == nil {
return 0
}
return int(*service.Spec.Mode.Replicated.Replicas)
}
func queueScaleFactor(numQueue int, queuePerMin int) float64 {
if numQueue > queuePerMin {
queuePressure := float64(numQueue) / float64(queuePerMin)
return 1.0 + math.Min(queuePressure-1.0, 1.0) return 1.0 + math.Min(queuePressure-1.0, 1.0)
} }
return 1.0 return 1.0
} }
func checkMemcached(dockercli *dockerclient.Client) (bool, error) { func checkMemcached(ctx context.Context, dockercli *dockerclient.Client) (bool, error) {
containerName := "shuffle-cache" containerName := "shuffle-cache"
continer, err := dockercli.ContainerInspect(context.Background(), containerName) continer, err := dockercli.ContainerInspect(context.Background(), containerName)
if err != nil { if err != nil {
@@ -3952,6 +3924,17 @@ func checkMemcached(dockercli *dockerclient.Client) (bool, error) {
return false, err return false, err
} }
if continer.State.Running == false {
log.Printf("[INFO] Container %s exists but is not running. Attempting to start it.", containerName)
err = dockercli.ContainerStart(ctx, containerName, container.StartOptions{})
if err != nil {
log.Printf("[ERROR] Failed to start container %s: %v", containerName, err)
return false, err
}
log.Printf("[INFO] Successfully started container %s.", containerName)
return true, nil
}
return continer.State.Running, nil return continer.State.Running, nil
} }
@@ -3970,12 +3953,6 @@ func deployMemcached(dockercli *dockerclient.Client) error {
Cmd: []string{"-m", defaultMem}, Cmd: []string{"-m", defaultMem},
} }
// PortBindings: nat.PortMap{
// "514/tcp": []nat.PortBinding{{HostPort: "514"}},
// "514/udp": []nat.PortBinding{{HostPort: "514"}},
// "5160/tcp": []nat.PortBinding{{HostPort: "5160"}},
// },
hostConfig := &container.HostConfig{ hostConfig := &container.HostConfig{
PortBindings: nat.PortMap{ PortBindings: nat.PortMap{
"11211/tcp": []nat.PortBinding{{HostPort: "11211"}}, "11211/tcp": []nat.PortBinding{{HostPort: "11211"}},
@@ -4015,6 +3992,7 @@ func nodesResourceUsage(ctx context.Context, client *dockerclient.Client) error
} }
*/ */
/*
func numberOfReplicas(ctx context.Context, queueLength int, config shuffle.ScalingConfig) (int, int) { func numberOfReplicas(ctx context.Context, queueLength int, config shuffle.ScalingConfig) (int, int) {
queueScaleFactor := queueScaleFactor(queueLength, config) queueScaleFactor := queueScaleFactor(queueLength, config)
numReplicas := int(float64(queueLength) * queueScaleFactor) numReplicas := int(float64(queueLength) * queueScaleFactor)
@@ -4052,7 +4030,10 @@ func numberOfReplicas(ctx context.Context, queueLength int, config shuffle.Scali
return numReplicas, runningReplicas return numReplicas, runningReplicas
} }
*/
// TODO: Currently we use number of request made for the worker to run a execution as it is much
// easier to track in a window time frame. But this could be useful.
func collectMetrics(ctx context.Context, dockerClient *dockerclient.Client) (int, error) { func collectMetrics(ctx context.Context, dockerClient *dockerclient.Client) (int, error) {
client := shuffle.GetExternalClient(baseUrl) client := shuffle.GetExternalClient(baseUrl)
fullUrl := fmt.Sprintf("%s/api/v1/workflows/queue", baseUrl) fullUrl := fmt.Sprintf("%s/api/v1/workflows/queue", baseUrl)