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]