This is an automated email from the ASF dual-hosted git repository.

pbacsko pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/yunikorn-core.git


The following commit(s) were added to refs/heads/master by this push:
     new 99ee3aa4 [YUNIKORN-2542] Consistent logging and tracker handling for 
increment/decrement (#876)
99ee3aa4 is described below

commit 99ee3aa469497853c42daf77bf0971ee00da3132
Author: Tseng Hsi-Huang <[email protected]>
AuthorDate: Fri May 31 08:37:37 2024 +0200

    [YUNIKORN-2542] Consistent logging and tracker handling for 
increment/decrement (#876)
    
    Closes: #876
    
    Signed-off-by: Peter Bacsko <[email protected]>
---
 pkg/scheduler/ugm/manager.go      | 31 +++++++++++++++++++++++++------
 pkg/scheduler/ugm/manager_test.go | 38 ++++++++++++++++++++++++++++++++++++++
 pkg/scheduler/ugm/user_tracker.go | 10 ----------
 3 files changed, 63 insertions(+), 16 deletions(-)

diff --git a/pkg/scheduler/ugm/manager.go b/pkg/scheduler/ugm/manager.go
index b374edbe..27044d52 100644
--- a/pkg/scheduler/ugm/manager.go
+++ b/pkg/scheduler/ugm/manager.go
@@ -97,8 +97,27 @@ func (m *Manager) IncreaseTrackedResource(queuePath, 
applicationID string, usage
        if !userTracker.hasGroupForApp(applicationID) {
                m.ensureGroupTrackerForApp(queuePath, applicationID, user)
        }
-
        userTracker.increaseTrackedResource(queuePath, applicationID, usage)
+       appGroup := userTracker.getGroupForApp(applicationID)
+       log.Log(log.SchedUGM).Debug("Increasing resource usage for user",
+               zap.String("user", user.User),
+               zap.String("queue path", queuePath),
+               zap.String("application", applicationID),
+               zap.String("group", appGroup),
+               zap.Stringer("resource", usage))
+       groupTracker := m.GetGroupTracker(appGroup)
+       if groupTracker == nil {
+               log.Log(log.SchedUGM).Error("group tracker should be available 
in groupTrackers map",
+                       zap.String("application", applicationID),
+                       zap.String("group", appGroup))
+               return
+       }
+       log.Log(log.SchedUGM).Debug("Increasing resource usage for group",
+               zap.String("group", appGroup),
+               zap.String("queue path", queuePath),
+               zap.String("application", applicationID),
+               zap.Stringer("resource", usage))
+       groupTracker.increaseTrackedResource(queuePath, applicationID, usage, 
userTracker.userName)
 }
 
 // DecreaseTrackedResource Decrease the resource usage for the given user 
group and queue path combination.
@@ -130,7 +149,7 @@ func (m *Manager) DecreaseTrackedResource(queuePath, 
applicationID string, usage
                zap.String("user", user.User),
                zap.String("queue path", queuePath),
                zap.String("application", applicationID),
-               zap.String("tracked group", appGroup),
+               zap.String("group", appGroup),
                zap.Stringer("resource", usage),
                zap.Bool("removeApp", removeApp))
        if userTracker.decreaseTrackedResource(queuePath, applicationID, usage, 
removeApp) {
@@ -145,8 +164,8 @@ func (m *Manager) DecreaseTrackedResource(queuePath, 
applicationID string, usage
        groupTracker := m.GetGroupTracker(appGroup)
        if groupTracker == nil {
                log.Log(log.SchedUGM).Error("group tracker should be available 
in groupTrackers map",
-                       zap.String("applicationID", applicationID),
-                       zap.String("applicationID", appGroup))
+                       zap.String("application", applicationID),
+                       zap.String("group", appGroup))
                return
        }
        log.Log(log.SchedUGM).Debug("Decreasing resource usage for group",
@@ -217,7 +236,7 @@ func (m *Manager) ensureGroupTrackerForApp(queuePath, 
applicationID string, user
                if groupTracker == nil {
                        log.Log(log.SchedUGM).Info("Group tracker doesn't 
exists. Creating appGroup tracker",
                                zap.String("queue path", queuePath),
-                               zap.String("appGroup", appGroup))
+                               zap.String("group", appGroup))
                        groupTracker = newGroupTracker(appGroup, m.events)
                        m.Lock()
                        m.groupTrackers[appGroup] = groupTracker
@@ -225,7 +244,7 @@ func (m *Manager) ensureGroupTrackerForApp(queuePath, 
applicationID string, user
                }
        }
        log.Log(log.SchedUGM).Info("Group tracker set for user application",
-               zap.String("appGroup", appGroup),
+               zap.String("group", appGroup),
                zap.String("user", user.User),
                zap.String("application", applicationID),
                zap.String("queue path", queuePath))
diff --git a/pkg/scheduler/ugm/manager_test.go 
b/pkg/scheduler/ugm/manager_test.go
index 81da27f7..9c571be9 100644
--- a/pkg/scheduler/ugm/manager_test.go
+++ b/pkg/scheduler/ugm/manager_test.go
@@ -719,6 +719,44 @@ func TestDecreaseTrackedResourceForGroupTracker(t 
*testing.T) {
        assert.Equal(t, 
resources.Equals(groupTracker.queueTracker.childQueueTrackers["parent"].resourceUsage,
 resources.Zero), true)
 }
 
+func TestIncreaseTrackedResourceForGroupTracker(t *testing.T) {
+       setupUGM()
+       // Queue setup:
+       // root->parent
+       user := security.UserGroup{User: "user1", Groups: []string{"group1"}}
+       conf := createConfigWithoutLimits()
+       conf.Queues[0].Queues[0].Limits = []configs.Limit{
+               {
+                       Limit:           "parent queue limit for a specific 
group",
+                       Groups:          user.Groups,
+                       MaxResources:    map[string]string{"memory": "100"},
+                       MaxApplications: 2,
+               },
+       }
+       manager := GetUserManager()
+       assert.NilError(t, manager.UpdateConfig(conf.Queues[0], "root"))
+
+       usage1, err := 
resources.NewResourceFromConf(map[string]string{"memory": "50"})
+       if err != nil {
+               t.Errorf("new resource create returned error or wrong resource: 
error %t, res %v", err, usage1)
+       }
+
+       manager.IncreaseTrackedResource("root.parent", TestApp1, usage1, user)
+       groupTracker := m.GetGroupTracker(user.Groups[0])
+       assert.Equal(t, groupTracker != nil, true)
+       assert.Equal(t, 
groupTracker.queueTracker.childQueueTrackers["parent"].runningApplications[TestApp1],
 true)
+       assert.Equal(t, 
resources.Equals(groupTracker.queueTracker.childQueueTrackers["parent"].resourceUsage,
 usage1), true)
+
+       usage2, err := 
resources.NewResourceFromConf(map[string]string{"memory": "30"})
+       if err != nil {
+               t.Errorf("new resource create returned error or wrong resource: 
error %t, res %v", err, usage2)
+       }
+
+       manager.IncreaseTrackedResource("root.parent", TestApp2, usage2, user)
+       assert.Equal(t, 
groupTracker.queueTracker.childQueueTrackers["parent"].runningApplications[TestApp2],
 true)
+       assert.Equal(t, 
resources.Equals(groupTracker.queueTracker.childQueueTrackers["parent"].resourceUsage,
 resources.Add(usage1, usage2)), true)
+}
+
 func TestUserGroupLimitWithMultipleApps(t *testing.T) {
        // Increase app rsources to different child queue, which have different 
group linkage
        // Queue setup:
diff --git a/pkg/scheduler/ugm/user_tracker.go 
b/pkg/scheduler/ugm/user_tracker.go
index f06581b3..56996464 100644
--- a/pkg/scheduler/ugm/user_tracker.go
+++ b/pkg/scheduler/ugm/user_tracker.go
@@ -21,13 +21,10 @@ package ugm
 import (
        "strings"
 
-       "go.uber.org/zap"
-
        "github.com/apache/yunikorn-core/pkg/common"
        "github.com/apache/yunikorn-core/pkg/common/configs"
        "github.com/apache/yunikorn-core/pkg/common/resources"
        "github.com/apache/yunikorn-core/pkg/locking"
-       "github.com/apache/yunikorn-core/pkg/log"
        "github.com/apache/yunikorn-core/pkg/webservice/dao"
 )
 
@@ -62,13 +59,6 @@ func (ut *UserTracker) increaseTrackedResource(queuePath 
string, applicationID s
        ut.events.sendIncResourceUsageForUser(ut.userName, queuePath, usage)
        hierarchy := strings.Split(queuePath, configs.DOT)
        ut.queueTracker.increaseTrackedResource(hierarchy, applicationID, user, 
usage)
-       gt := ut.appGroupTrackers[applicationID]
-       log.Log(log.SchedUGM).Debug("Increasing resource usage for group",
-               zap.String("group", gt.getName()),
-               zap.Strings("queue path", hierarchy),
-               zap.String("application", applicationID),
-               zap.Stringer("resource", usage))
-       gt.increaseTrackedResource(queuePath, applicationID, usage, ut.userName)
 }
 
 func (ut *UserTracker) decreaseTrackedResource(queuePath string, applicationID 
string, usage *resources.Resource, removeApp bool) bool {


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to