This is an automated email from the ASF dual-hosted git repository.
ccondit 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 5d04f9c0 [YUNIKORN-1797] Core: Implement scoped logging (#566)
5d04f9c0 is described below
commit 5d04f9c0dd32d4c74cdc301cdf816807566d12e5
Author: Craig Condit <[email protected]>
AuthorDate: Tue Jun 20 09:11:12 2023 -0500
[YUNIKORN-1797] Core: Implement scoped logging (#566)
Add API support for scoped logging to the scheduler core. This will allow
for
fine-grained logging control per-subsystem. Also deprecates and ignores REST
API calls to get / reset logging level, as this is now controlled per
logger.
Closes: #566
Signed-off-by: Craig Condit <[email protected]>
---
Makefile | 6 +
pkg/common/configs/configs.go | 7 +
pkg/log/filtered_core.go | 54 ++++++
pkg/log/logger.go | 284 ++++++++++++++++++++++++------
pkg/log/logger_test.go | 294 +++++++++++++++++++++++++-------
pkg/scheduler/objects/node.go | 3 +-
pkg/scheduler/placement/filter.go | 42 ++---
pkg/scheduler/placement/placement.go | 10 +-
pkg/scheduler/tests/performance_test.go | 8 +-
pkg/webservice/handlers.go | 15 +-
pkg/webservice/handlers_test.go | 17 +-
pkg/webservice/state_dump.go | 2 +-
12 files changed, 573 insertions(+), 169 deletions(-)
diff --git a/Makefile b/Makefile
index 48ed1f37..fae7bb9a 100644
--- a/Makefile
+++ b/Makefile
@@ -160,6 +160,12 @@ test: clean
go test ./... $(RACE) -tags deadlock -coverprofile=coverage.txt
-covermode=atomic
go vet $(REPO)...
+# Run benchmarks
+.PHONY: bench
+bench:
+ @echo "running benchmarks"
+ go test -v -run '^Benchmark' -bench . ./pkg/...
+
# Generate FSM graphs (dot/png)
.PHONY: fsm_graph
fsm_graph: clean
diff --git a/pkg/common/configs/configs.go b/pkg/common/configs/configs.go
index 52803742..201a8d76 100644
--- a/pkg/common/configs/configs.go
+++ b/pkg/common/configs/configs.go
@@ -21,6 +21,8 @@ package configs
import (
"sync"
"time"
+
+ "github.com/apache/yunikorn-core/pkg/log"
)
const (
@@ -42,6 +44,11 @@ func init() {
configs: make(map[string]*SchedulerConfig),
lock: &sync.RWMutex{},
}
+
+ // add a callback to reconfigure logging
+ AddConfigMapCallback("logging", func() {
+ log.UpdateLoggingConfig(GetConfigMap())
+ })
}
// scheduler config context provides thread-safe access for scheduler
configurations
diff --git a/pkg/log/filtered_core.go b/pkg/log/filtered_core.go
new file mode 100644
index 00000000..1e2dcf6a
--- /dev/null
+++ b/pkg/log/filtered_core.go
@@ -0,0 +1,54 @@
+/*
+ Licensed to the Apache Software Foundation (ASF) under one
+ or more contributor license agreements. See the NOTICE file
+ distributed with this work for additional information
+ regarding copyright ownership. The ASF licenses this file
+ to you under the Apache License, Version 2.0 (the
+ "License"); you may not use this file except in compliance
+ with the License. You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+ Unless required by applicable law or agreed to in writing, software
+ distributed under the License is distributed on an "AS IS" BASIS,
+ WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ See the License for the specific language governing permissions and
+ limitations under the License.
+*/
+
+package log
+
+import "go.uber.org/zap/zapcore"
+
+type filteredCore struct {
+ level zapcore.Level
+ inner zapcore.Core
+}
+
+var _ zapcore.Core = filteredCore{}
+
+func (f filteredCore) Enabled(level zapcore.Level) bool {
+ if level < f.level {
+ return false
+ }
+ return f.inner.Enabled(level)
+}
+
+func (f filteredCore) With(fields []zapcore.Field) zapcore.Core {
+ return f.inner.With(fields)
+}
+
+func (f filteredCore) Check(entry zapcore.Entry, ce *zapcore.CheckedEntry)
*zapcore.CheckedEntry {
+ if entry.Level < f.level {
+ return ce
+ }
+ return f.inner.Check(entry, ce)
+}
+
+func (f filteredCore) Write(entry zapcore.Entry, fields []zapcore.Field) error
{
+ return f.inner.Write(entry, fields)
+}
+
+func (f filteredCore) Sync() error {
+ return f.inner.Sync()
+}
diff --git a/pkg/log/logger.go b/pkg/log/logger.go
index b56fda18..e3b52287 100644
--- a/pkg/log/logger.go
+++ b/pkg/log/logger.go
@@ -19,10 +19,12 @@
package log
import (
- "errors"
"fmt"
"reflect"
+ "strconv"
+ "strings"
"sync"
+ "sync/atomic"
"go.uber.org/zap"
"go.uber.org/zap/zapcore"
@@ -30,45 +32,149 @@ import (
var once sync.Once
var logger *zap.Logger
-var config *zap.Config
+var zapConfigs *zap.Config
+// LoggerHandle is used to efficiently look up logger references
+type LoggerHandle struct {
+ id int
+ name string
+}
+
+func (h LoggerHandle) String() string {
+ return h.name
+}
+
+// Logger constants for configuration
+const (
+ nullLogger = ""
+ defaultLog = "log.level"
+ logPrefix = "log."
+ levelSuffix = ".level"
+)
+
+// Predefined loggers: when adding new loggers, ids must be sequential, and
all must be added to the loggers slice in the same order
+var (
+ Core = &LoggerHandle{id: 1, name: "core"}
+ Test = &LoggerHandle{id: 2, name: "test"}
+)
+
+// this tracks all the known logger handles, used to preallocate the real
logger instances when configuration changes
+var loggers = []*LoggerHandle{
+ Core,
+ Test,
+}
+
+// structure to hold all current logger configuration state
+type loggerConfig struct {
+ loggers []*zap.Logger
+}
+
+// tracks the currently used set of loggers; replaced completely whenever
configuration changes
+var currentLoggerConfig = atomic.Pointer[loggerConfig]{}
+
+// tracks the default logger handle, which is used for legacy log.Logger()
calls
+var defaultLogger = atomic.Pointer[LoggerHandle]{}
+
+// Logger retrieves the global logger. This is for compatibility with legacy
code and
+// should not be used for new log messages; use Log(loggerHandle) instead
func Logger() *zap.Logger {
- once.Do(func() {
- if logger = zap.L(); isNopLogger(logger) {
- // If a global logger is not found, this could be
either scheduler-core
- // is running as a deployment mode, or running with
another non-go code
- // shim. In this case, we need to create our own logger.
- // TODO support log options when a global logger is not
there
- config = createConfig()
- var err error
- logger, err = config.Build()
- // this should really not happen so just write to
stdout and set a Nop logger
- if err != nil {
- fmt.Printf("Logging disabled, logger init
failed with error: %v\n", err)
- logger = zap.NewNop()
- }
- }
- })
+ once.Do(initLogger)
+ return Log(defaultLogger.Load())
+}
+// RootLogger retrieves the root logger, visible for testing
+func RootLogger() *zap.Logger {
+ once.Do(initLogger)
return logger
}
+// Log retrieves a named logger
+func Log(handle *LoggerHandle) *zap.Logger {
+ once.Do(initLogger)
+ if handle == nil || handle.id == 0 {
+ handle = defaultLogger.Load()
+ }
+ conf := currentLoggerConfig.Load()
+ return conf.loggers[handle.id-1]
+}
+
+// createLogger creates a new, named logger that has log levels filtered based
on the given configuration.
+// levelMap contains a mapping of fully-qualified logger names to log levels
+func createLogger(levelMap map[string]zapcore.Level, name string) *zap.Logger {
+ level := loggerLevel(levelMap, name)
+ return logger.Named(name).WithOptions(zap.WrapCore(func(inner
zapcore.Core) zapcore.Core {
+ return filteredCore{inner: inner, level: level}
+ }))
+}
+
+// loggerLevel computes the log level that should be associated with an
arbitrarily named logger.
+// The levelMap is used to look up the level. If not found, the parent logger
is looked up until no further
+// parents can be tried, consulting levelMap for each possible match. For
example, given a logger named
+// core.context.cache, the following keys will be looked up in levelMap:
+//
+// "core.context.cache"
+// "core.context"
+// "core"
+// "" (default logger)
+//
+// The first key that returns a match determines the log level returned. This
allows a level of
+// "core" to control all the child loggers as well if there is not a more
specific override.
+// If no configuration can be found, even for the default empty logger,
InfoLevel will be returned.
+func loggerLevel(levelMap map[string]zapcore.Level, name string) zapcore.Level
{
+ for ; name != nullLogger; name = parentLogger(name) {
+ if level, ok := levelMap[name]; ok {
+ return level
+ }
+ }
+ if level, ok := levelMap[nullLogger]; ok {
+ return level
+ }
+ return zapcore.InfoLevel
+}
+
+// parentLogger returns the name of the parent logger or the empty string ""
if no parent exists.
+// Loggers are named with periods (.) as separators; i.e. the parent logger of
"a.b.c" is "a.b", the parent of
+// "a.b" is "a", and the parent of "a" is "".
+func parentLogger(name string) string {
+ i := strings.LastIndex(name, ".")
+ if i < 0 {
+ // no remaining periods; parent is the null logger
+ return nullLogger
+ }
+ return name[0:i]
+}
+
+func initLogger() {
+ if logger = zap.L(); isNopLogger(logger) {
+ // If a global logger is not found, this could be either
scheduler-core
+ // is running as a deployment mode, or running with another
non-go code
+ // shim. In this case, we need to create our own logger.
+ zapConfigs = createConfig()
+ var err error
+ logger, err = zapConfigs.Build()
+ // this should really not happen so just write to stdout and
set a Nop logger
+ if err != nil {
+ fmt.Printf("Logging disabled, logger init failed with
error: %v\n", err)
+ logger = zap.NewNop()
+ }
+ }
+
+ // set default logger
+ defaultLogger.Store(Core)
+
+ // initialize sub-loggers
+ initLoggingConfig(nil)
+}
+
func InitializeLogger(log *zap.Logger, zapConfig *zap.Config) {
once.Do(func() {
logger = log
- config = zapConfig
- logger.Info("Using an already initialized logger")
+ zapConfigs = zapConfig
+ defaultLogger.Store(Core)
+ initLoggingConfig(nil)
})
}
-func IsDebugEnabled() bool {
- if logger == nil {
- // when under development mode
- return true
- }
- return logger.Core().Enabled(zapcore.DebugLevel)
-}
-
// Returns true if the logger is a noop.
// Logger is a noop means the logger has not been initialized yet.
// This usually means a global logger is not set in the given context,
@@ -78,14 +184,6 @@ func isNopLogger(logger *zap.Logger) bool {
return reflect.DeepEqual(zap.NewNop(), logger)
}
-// Visible by tests
-func InitAndSetLevel(level zapcore.Level) {
- if config == nil {
- Logger()
- }
- config.Level.SetLevel(level)
-}
-
// Create a log config to keep full control over
// LogLevel set to DEBUG, Encodes for console, Writes to stderr,
// Enables development mode (DPanicLevel),
@@ -99,7 +197,7 @@ func createConfig() *zap.Config {
MessageKey: "message",
LevelKey: "level",
TimeKey: "time",
- NameKey: "name",
+ NameKey: "logger",
CallerKey: "caller",
StacktraceKey: "stacktrace",
LineEnding: zapcore.DefaultLineEnding,
@@ -115,27 +213,105 @@ func createConfig() *zap.Config {
}
}
-func SetLogLevel(newLevel string) error {
- oldLevel := config.Level.String()
+func GetZapConfigs() *zap.Config {
+ return zapConfigs
+}
+
+// UpdateLoggingConfig is used to reconfigure logging. This uses config keys
of the form log.{logger}.level={level}.
+// The default level is set by log.level={level}. The {level} value can be
either numeric (-1 through 5), or
+// textual (DEBUG, INFO, WARN, ERROR, DPANIC, PANIC, or ERROR). See zapcore
documentation for more details.
+func UpdateLoggingConfig(config map[string]string) {
+ once.Do(initLogger)
+ initLoggingConfig(config)
+}
+
+// initLoggingConfig replaces the existing set of loggers with new ones
configured according to the given
+// configuration. All keys of the form "log.{name}.level" will be parsed and a
map of logger name -> logger level
+// will be created. For each defined logger handle (see above), a new logger
instance is created using the
+// most specific configuration found. For example, a logger named "a.b.c" will
be configured by "log.a.b.c.level".
+// If this key does not exist, "a.b" will be consulted, then "a", and finally
"". If no configuration is found for a
+// given key, Info will be used. Finally, the finest log level specified in
any configuration will be used to set the
+// zap log level of the root logger. So if two loggers (one INFO, and one
DEBUG) are configured, the root logger will
+// be set to DEBUG level. If both loggers were at INFO level, the root logger
would be set to INFO.
+// Each configured logger will filter log messages at its own level by
wrapping the underlying zap.Core implementation
+// with one that checks the enabled log level first.
+func initLoggingConfig(config map[string]string) {
+ levelMap := make(map[string]zapcore.Level)
+ levelMap[nullLogger] = zapcore.InfoLevel
+ zapLoggers := make([]*zap.Logger, len(loggers))
- // noop if the input is the same as what is set
- if newLevel == oldLevel {
- return nil
+ // override default level if found (log.level key)
+ if defaultLevel, ok := config[defaultLog]; ok {
+ if levelRef := parseLevel(defaultLevel); levelRef != nil {
+ levelMap[nullLogger] = *levelRef
+ }
}
- logger.Info("Updating log level",
- zap.String("new level", newLevel))
- text := []byte(newLevel)
- if err := config.Level.UnmarshalText(text); err != nil {
- var errorMsg = "failed to change log level, old level active"
- logger.Error(errorMsg, zap.String("loglevel", oldLevel))
- return errors.New(errorMsg)
+ // parse out log entries and build level map
+ for k, v := range config {
+ // disallow spaces and double periods
+ if strings.Contains(k, "..") || strings.Contains(k, " ") {
+ continue
+ }
+ // ensure config key starts with "log."
+ name, ok := strings.CutPrefix(k, logPrefix)
+ if !ok {
+ continue
+ }
+ // / ensure config key ends with ".level"
+ name, ok = strings.CutSuffix(name, levelSuffix)
+ if !ok {
+ continue
+ }
+ // if level is a valid log level, store it in the level map
+ if levelRef := parseLevel(v); levelRef != nil {
+ levelMap[name] = *levelRef
+ }
}
- logger.Info("Log level updated", zap.String("old level", oldLevel),
- zap.String("new level", newLevel))
- return nil
+
+ // compute the finest log level necessary to allow all loggers to
succeed
+ minLevel := zapcore.InvalidLevel - 1
+ for _, v := range levelMap {
+ if minLevel > v {
+ minLevel = v
+ }
+ }
+
+ // create each configured logger and initialize the overall
configuration
+ for i := 0; i < len(loggers); i++ {
+ zapLoggers[i] = createLogger(levelMap, loggers[i].name)
+ }
+ newLoggerConfig := loggerConfig{loggers: zapLoggers}
+
+ // update the root zap logger level
+ zapConfigs.Level.SetLevel(minLevel)
+
+ // atomically update the set of loggers
+ currentLoggerConfig.Store(&newLoggerConfig)
}
-func GetConfig() *zap.Config {
- return config
+// parseLevel parses a textual (or numeric) log level into a zapcore.Level
instance. Both numeric (-1 <= level <= 5)
+// and textual (DEBUG, INFO, WARN, ERROR, DPANIC, PANIC, FATAL) are supported.
+func parseLevel(level string) *zapcore.Level {
+ // parse text
+ zapLevel, err := zapcore.ParseLevel(level)
+ if err == nil {
+ return &zapLevel
+ }
+
+ // parse numeric
+ levelNum, err := strconv.ParseInt(level, 10, 31)
+ if err == nil {
+ zapLevel = zapcore.Level(levelNum)
+ if zapLevel < zapcore.DebugLevel {
+ zapLevel = zapcore.DebugLevel
+ }
+ if zapLevel >= zapcore.InvalidLevel {
+ zapLevel = zapcore.InvalidLevel - 1
+ }
+ return &zapLevel
+ }
+
+ // parse failed
+ return nil
}
diff --git a/pkg/log/logger_test.go b/pkg/log/logger_test.go
index 264f14de..1e25037c 100644
--- a/pkg/log/logger_test.go
+++ b/pkg/log/logger_test.go
@@ -19,14 +19,36 @@
package log
import (
+ "fmt"
+ "os"
"sync"
"testing"
+ "time"
"go.uber.org/zap"
"go.uber.org/zap/zapcore"
"gotest.tools/v3/assert"
)
+var logDir string
+var logFile string
+
+var iterations = 100000
+
+func TestLoggerIds(t *testing.T) {
+ _ = Logger()
+
+ // validate logger count
+ assert.Equal(t, 2, len(loggers), "wrong logger count")
+
+ // validate that all loggers are populated and have sequential ids
+ for i := 0; i < len(loggers); i++ {
+ handle := loggers[i]
+ assert.Assert(t, handle != nil, "nil handle for index", i)
+ assert.Equal(t, handle.id, i+1, "wrong id", handle.name)
+ }
+}
+
// This test sets the global zap logger. This must be undone to make sure no
side
// effects on other tests are caused by running this test.
func TestIsNopLogger(t *testing.T) {
@@ -50,35 +72,10 @@ func TestIsNopLogger(t *testing.T) {
assert.Equal(t, false, isNopLogger(zap.L()))
}
-// Since we test the function IsDebugEnabled() we set the logger global var.
-// It has not triggered the once.Do() so we just need to make sure we clean up
the
-// global var.
-func TestIsDebugEnabled(t *testing.T) {
- // reset the global vars and zap logger
- defer resetGlobals()
-
- zapConfig := zap.Config{
- Level: zap.NewAtomicLevelAt(zapcore.DebugLevel),
- Encoding: "console",
- }
- var err error
- logger, err = zapConfig.Build()
- assert.NilError(t, err, "debug level logger create failed")
- assert.Equal(t, true, IsDebugEnabled())
-
- zapConfig = zap.Config{
- Level: zap.NewAtomicLevelAt(zapcore.InfoLevel),
- Encoding: "console",
- }
- logger, err = zapConfig.Build()
- assert.NilError(t, err, "info level logger create failed")
- assert.Equal(t, false, IsDebugEnabled())
-}
-
// reset the global vars and the global logger in zap
func resetGlobals() {
logger = nil
- config = nil
+ zapConfigs = nil
once = sync.Once{}
zap.ReplaceGlobals(zap.NewNop())
}
@@ -98,11 +95,6 @@ func TestCreateConfig(t *testing.T) {
assert.Assert(t, logger == nil, "global logger should not have been set
%v", logger)
localLogger = Logger()
assert.Assert(t, localLogger != nil, "returned logger should have been
not nil")
- // default log level is debug
- assert.Equal(t, true, IsDebugEnabled())
- // change log level to info
- InitAndSetLevel(zap.InfoLevel)
- assert.Equal(t, false, IsDebugEnabled())
}
func TestInitializeLogger(t *testing.T) {
@@ -118,44 +110,230 @@ func TestInitializeLogger(t *testing.T) {
assert.NilError(t, err2, "failed to create local logger")
InitializeLogger(localLogger, &zapConfig)
- assert.Equal(t, Logger(), localLogger)
+ assert.Equal(t, RootLogger(), localLogger)
// second initialization should not do anything
InitializeLogger(localLogger2, &zapConfig)
- assert.Equal(t, Logger(), localLogger)
+ assert.Equal(t, RootLogger(), localLogger)
+}
+func BenchmarkLegacyLoggerDebug(b *testing.B) {
+ benchmarkLegacyLoggerDebug(b.N)
}
-func TestChangeValidLogLevel(t *testing.T) {
- defer resetGlobals()
+func TestLegacyLoggerDebug(t *testing.T) {
+ nsOp := benchmarkLegacyLoggerDebug(iterations)
+ RootLogger().Info("log.Logger() performance", zap.Int64("debug
(ns/op)", nsOp))
+}
- zapConfig := zap.Config{
- Level: zap.NewAtomicLevelAt(zapcore.InfoLevel),
- Encoding: "console",
+func benchmarkLegacyLoggerDebug(iterations int) int64 {
+ _ = Logger()
+ initTestLogger()
+ defer resetTestLogger()
+ start := time.Now()
+ for i := 0; i < iterations; i++ {
+ RootLogger().Debug("test", zap.String("foo", "bar"))
}
- localLogger, err := zapConfig.Build()
- assert.NilError(t, err, "failed to create local logger")
- InitializeLogger(localLogger, &zapConfig)
+ return (time.Since(start).Nanoseconds()) / int64(iterations)
+}
+
+func BenchmarkLegacyLoggerInfo(b *testing.B) {
+ benchmarkLegacyLoggerInfo(b.N)
+}
- err = SetLogLevel("DEBUG")
- assert.NilError(t, err, "failed to change log level")
- assert.Equal(t, zapConfig.Level.Level(), zapcore.DebugLevel)
+func TestLegacyLoggerInfo(t *testing.T) {
+ nsOp := benchmarkLegacyLoggerInfo(iterations)
+ RootLogger().Info("log.Logger() performance", zap.Int64("info (ns/op)",
nsOp))
+}
- // set again to see that we keep DEBUG without issues
- err = SetLogLevel("DEBUG")
- assert.NilError(t, err, "failed to change log level")
- assert.Equal(t, zapConfig.Level.Level(), zapcore.DebugLevel)
+func benchmarkLegacyLoggerInfo(iterations int) int64 {
+ _ = Logger()
+ initTestLogger()
+ defer resetTestLogger()
+ start := time.Now()
+ for i := 0; i < iterations; i++ {
+ RootLogger().Info("test", zap.String("foo", "bar"))
+ }
+ return (time.Since(start).Nanoseconds()) / int64(iterations)
}
-func TestChangeInvalidLogLevel(t *testing.T) {
- defer resetGlobals()
+func BenchmarkScopedLoggerDebug(b *testing.B) {
+ benchmarkScopedLoggerDebug(b.N)
+}
- zapConfig := zap.Config{
- Level: zap.NewAtomicLevelAt(zapcore.InfoLevel),
- Encoding: "console",
+func TestScopedLoggerDebug(t *testing.T) {
+ nsOp := benchmarkScopedLoggerDebug(iterations)
+ Log(Test).Info("log.Log(...) performance (root=INFO)", zap.Int64("debug
(ns/op)", nsOp))
+}
+
+func benchmarkScopedLoggerDebug(iterations int) int64 {
+ _ = Logger()
+ initTestLogger()
+ defer resetTestLogger()
+ UpdateLoggingConfig(map[string]string{
+ "log.level": "INFO",
+ })
+ start := time.Now()
+ for i := 0; i < iterations; i++ {
+ Log(Core).Debug("test", zap.String("foo", "bar"))
}
- localLogger, err := zapConfig.Build()
- assert.NilError(t, err, "default config logger create failed")
- InitializeLogger(localLogger, &zapConfig)
+ return (time.Since(start).Nanoseconds()) / int64(iterations)
+}
+
+func BenchmarkScopedLoggerInfo(b *testing.B) {
+ benchmarkScopedLoggerInfo(b.N)
+}
+
+func TestScopedLoggerInfo(t *testing.T) {
+ nsOp := benchmarkScopedLoggerInfo(iterations)
+ Log(Test).Info("log.Log(...) performance (root=INFO)", zap.Int64("info
(ns/op)", nsOp))
+}
+
+func benchmarkScopedLoggerInfo(iterations int) int64 {
+ _ = Logger()
+ initTestLogger()
+ defer resetTestLogger()
+ UpdateLoggingConfig(map[string]string{
+ "log.level": "INFO",
+ })
+ start := time.Now()
+ for i := 0; i < iterations; i++ {
+ Log(Core).Info("test", zap.String("foo", "bar"))
+ }
+ return (time.Since(start).Nanoseconds()) / int64(iterations)
+}
+
+func BenchmarkScopedLoggerDebugEnabled(b *testing.B) {
+ benchmarkScopedLoggerDebugEnabled(b.N)
+}
+
+func TestScopedLoggerDebugEnabled(t *testing.T) {
+ nsOp := benchmarkScopedLoggerDebugEnabled(iterations)
+ Log(Test).Info("log.Log(...) performance (root=DEBUG)",
zap.Int64("debug (ns/op)", nsOp))
+}
+
+func benchmarkScopedLoggerDebugEnabled(iterations int) int64 {
+ _ = Logger()
+ initTestLogger()
+ defer resetTestLogger()
+ UpdateLoggingConfig(map[string]string{
+ "log.test.level": "DEBUG",
+ })
+ start := time.Now()
+ for i := 0; i < iterations; i++ {
+ Log(Test).Debug("test", zap.String("foo", "bar"))
+ }
+ return (time.Since(start).Nanoseconds()) / int64(iterations)
+}
+
+func BenchmarkScopedLoggerInfoFiltered(b *testing.B) {
+ benchmarkScopedLoggerInfoFiltered(b.N)
+}
- err = SetLogLevel("INVALID")
- assert.Error(t, err, "failed to change log level, old level active")
+func TestScopedLoggerInfoFiltered(t *testing.T) {
+ nsOp := benchmarkScopedLoggerInfoFiltered(iterations)
+ Log(Test).Info("log.Log(...) performance (root=DEBUG)", zap.Int64("info
(ns/op)", nsOp))
+}
+
+func benchmarkScopedLoggerInfoFiltered(iterations int) int64 {
+ _ = Logger()
+ initTestLogger()
+ defer resetTestLogger()
+ UpdateLoggingConfig(map[string]string{
+ "log.test.level": "DEBUG",
+ })
+ start := time.Now()
+ for i := 0; i < iterations; i++ {
+ Log(Core).Info("test", zap.String("foo", "bar"))
+ }
+ return (time.Since(start).Nanoseconds()) / int64(iterations)
+}
+
+func TestParseLevel(t *testing.T) {
+ assert.Equal(t, zapcore.DebugLevel, *parseLevel("-2"), "out of range
low")
+ assert.Equal(t, zapcore.DebugLevel, *parseLevel("-1"))
+ assert.Equal(t, zapcore.InfoLevel, *parseLevel("0"))
+ assert.Equal(t, zapcore.WarnLevel, *parseLevel("1"))
+ assert.Equal(t, zapcore.ErrorLevel, *parseLevel("2"))
+ assert.Equal(t, zapcore.DPanicLevel, *parseLevel("3"))
+ assert.Equal(t, zapcore.PanicLevel, *parseLevel("4"))
+ assert.Equal(t, zapcore.FatalLevel, *parseLevel("5"))
+ assert.Equal(t, zapcore.FatalLevel, *parseLevel("6"), "out of range
high")
+ assert.Assert(t, parseLevel("+2-3") == nil, "parse error")
+ assert.Equal(t, zapcore.DebugLevel, *parseLevel("Debug"))
+ assert.Equal(t, zapcore.InfoLevel, *parseLevel("iNFO"))
+ assert.Equal(t, zapcore.WarnLevel, *parseLevel("WaRn"))
+ assert.Equal(t, zapcore.ErrorLevel, *parseLevel("ERROR"))
+ assert.Equal(t, zapcore.DPanicLevel, *parseLevel("dpanic"))
+ assert.Equal(t, zapcore.PanicLevel, *parseLevel("PAnIC"))
+ assert.Equal(t, zapcore.FatalLevel, *parseLevel("faTal"))
+ assert.Assert(t, parseLevel("x") == nil, "parse error")
+}
+
+func TestParentLogger(t *testing.T) {
+ assert.Equal(t, "", parentLogger(""), "nullLogger")
+ assert.Equal(t, "", parentLogger("a"), "level 1")
+ assert.Equal(t, "a", parentLogger("a.b"), "level 2")
+ assert.Equal(t, "a.b", parentLogger("a.b.c"), "level 3")
+}
+
+func resetTestLogger() {
+ // flush log
+ logger.Sync() //nolint:errcheck
+
+ // init default logger
+ initLogger()
+
+ // update logger config to defaults
+ UpdateLoggingConfig(map[string]string{})
+
+ if logFile != "" {
+ logFile = ""
+ }
+ if logDir != "" {
+ if err := os.RemoveAll(logDir); err != nil {
+ fmt.Printf("Error removing log dir: %s", err.Error())
+ }
+ }
+}
+
+// initTestLogger is basically the same as the default initLogger() function
but uses a temporary file.
+// this ensures that the logging API is actually used, while allowing us to
avoid massive log spam to stdout
+func initTestLogger() {
+ path, err := os.MkdirTemp("", "log*")
+ if err != nil {
+ panic(err)
+ }
+ logDir = path
+ logFile = fmt.Sprintf("%s/log.stdout", logDir)
+ outputPaths := []string{logFile}
+ zapConfigs = &zap.Config{
+ Level: zap.NewAtomicLevelAt(zapcore.Level(0)),
+ Development: false,
+ DisableCaller: false,
+ DisableStacktrace: false,
+ Sampling: nil,
+ Encoding: "console",
+ EncoderConfig: zapcore.EncoderConfig{
+ MessageKey: "message",
+ LevelKey: "level",
+ TimeKey: "time",
+ NameKey: "logger",
+ CallerKey: "caller",
+ StacktraceKey: "stacktrace",
+ LineEnding: zapcore.DefaultLineEnding,
+ // note:
https://godoc.org/go.uber.org/zap/zapcore#EncoderConfig
+ // only EncodeName is optional all others must be set
+ EncodeLevel: zapcore.CapitalLevelEncoder,
+ EncodeTime: zapcore.ISO8601TimeEncoder,
+ EncodeDuration: zapcore.StringDurationEncoder,
+ EncodeCaller: zapcore.ShortCallerEncoder,
+ },
+ OutputPaths: outputPaths,
+ ErrorOutputPaths: []string{"stderr"},
+ }
+
+ logger, err = zapConfigs.Build()
+ if err != nil {
+ panic(err)
+ }
+ defer logger.Sync() //nolint:errcheck
}
diff --git a/pkg/scheduler/objects/node.go b/pkg/scheduler/objects/node.go
index d1a68e40..6e57ec1e 100644
--- a/pkg/scheduler/objects/node.go
+++ b/pkg/scheduler/objects/node.go
@@ -70,7 +70,8 @@ func NewNode(proto *si.NodeInfo) *Node {
var ready bool
var err error
if ready, err =
strconv.ParseBool(proto.Attributes[common.NodeReadyAttribute]); err != nil {
- log.Logger().Error("Could not parse ready flag, assuming true")
+ log.Logger().Debug("Could not parse ready flag, assuming true",
+ zap.String("nodeID", proto.NodeID))
ready = true
}
sn := &Node{
diff --git a/pkg/scheduler/placement/filter.go
b/pkg/scheduler/placement/filter.go
index 47d52bd3..b2c296a9 100644
--- a/pkg/scheduler/placement/filter.go
+++ b/pkg/scheduler/placement/filter.go
@@ -168,25 +168,27 @@ func newFilter(conf configs.Filter) Filter {
}
// log the filter with all details (only at debug)
- if log.IsDebugEnabled() {
- var userfilter, groupfilter string
- if filter.userExp == nil {
- userfilter = "nil"
- } else {
- userfilter = filter.userExp.String()
- }
- if filter.groupExp == nil {
- groupfilter = "nil"
- } else {
- groupfilter = filter.groupExp.String()
- }
- log.Logger().Debug("Filter creation passed",
- zap.Bool("allow", filter.allow),
- zap.Bool("empty", filter.empty),
- zap.Any("userList", filter.userList),
- zap.Any("groupList", filter.groupList),
- zap.String("userFilter", userfilter),
- zap.String("groupFilter", groupfilter))
- }
+ logFilter(&filter)
return filter
}
+
+func logFilter(filter *Filter) {
+ var userfilter, groupfilter string
+ if filter.userExp == nil {
+ userfilter = "nil"
+ } else {
+ userfilter = filter.userExp.String()
+ }
+ if filter.groupExp == nil {
+ groupfilter = "nil"
+ } else {
+ groupfilter = filter.groupExp.String()
+ }
+ log.Logger().Debug("Filter creation passed",
+ zap.Bool("allow", filter.allow),
+ zap.Bool("empty", filter.empty),
+ zap.Any("userList", filter.userList),
+ zap.Any("groupList", filter.groupList),
+ zap.String("userFilter", userfilter),
+ zap.String("groupFilter", groupfilter))
+}
diff --git a/pkg/scheduler/placement/placement.go
b/pkg/scheduler/placement/placement.go
index a25524cc..f420c353 100644
--- a/pkg/scheduler/placement/placement.go
+++ b/pkg/scheduler/placement/placement.go
@@ -102,12 +102,10 @@ func (m *AppPlacementManager) initialise(rules
[]configs.PlacementRule) error {
m.rules = tempRules
// all done manager is initialised
m.initialised = true
- if log.IsDebugEnabled() {
- for rule := range m.rules {
- log.Logger().Debug("rule set",
- zap.Int("ruleNumber", rule),
- zap.String("ruleName", m.rules[rule].getName()))
- }
+ for rule := range m.rules {
+ log.Logger().Debug("rule set",
+ zap.Int("ruleNumber", rule),
+ zap.String("ruleName", m.rules[rule].getName()))
}
return nil
}
diff --git a/pkg/scheduler/tests/performance_test.go
b/pkg/scheduler/tests/performance_test.go
index 3bb0cf9d..4749b5c5 100644
--- a/pkg/scheduler/tests/performance_test.go
+++ b/pkg/scheduler/tests/performance_test.go
@@ -24,7 +24,6 @@ import (
"testing"
"time"
- "go.uber.org/zap"
"gotest.tools/v3/assert"
"github.com/apache/yunikorn-core/pkg/entrypoint"
@@ -33,7 +32,9 @@ import (
)
func benchmarkScheduling(b *testing.B, numNodes, numPods int) {
- log.InitAndSetLevel(zap.InfoLevel)
+ log.UpdateLoggingConfig(map[string]string{"log.level": "WARN"})
+ defer log.UpdateLoggingConfig(nil)
+
// Start all tests
serviceContext := entrypoint.StartAllServices()
defer serviceContext.StopAll()
@@ -71,6 +72,9 @@ partitions:
Version: "0.0.2",
BuildInfo: BuildInfoMap,
Config: configData,
+ ExtraConfig: map[string]string{
+ "log.level": "WARN",
+ },
}, mockRM)
assert.NilError(b, err, "RegisterResourceManager failed")
diff --git a/pkg/webservice/handlers.go b/pkg/webservice/handlers.go
index 2076dccb..1bdabc30 100644
--- a/pkg/webservice/handlers.go
+++ b/pkg/webservice/handlers.go
@@ -623,22 +623,13 @@ func getApplication(w http.ResponseWriter, r
*http.Request) {
func setLogLevel(w http.ResponseWriter, r *http.Request) {
writeHeaders(w)
- vars := httprouter.ParamsFromContext(r.Context())
- if vars == nil {
- buildJSONErrorResponse(w, MissingParamsName,
http.StatusBadRequest)
- return
- }
- level := vars.ByName("level")
- if err := log.SetLogLevel(level); err != nil {
- buildJSONErrorResponse(w, err.Error(), http.StatusBadRequest)
- }
+ log.Logger().Warn("Setting log levels via the REST API is deprecated.
The /ws/v1/loglevel endpoint will be removed in a future release.")
}
func getLogLevel(w http.ResponseWriter, r *http.Request) {
writeHeaders(w)
- zapConfig := log.GetConfig()
- if _, err := w.Write([]byte(zapConfig.Level.Level().String())); err !=
nil {
- log.Logger().Error("Could not get log level", zap.Error(err))
+ log.Logger().Warn("Getting log levels via the REST API is deprecated.
The /ws/v1/loglevel endpoint will be removed in a future release.")
+ if _, err := w.Write([]byte("info")); err != nil {
buildJSONErrorResponse(w, err.Error(),
http.StatusInternalServerError)
}
}
diff --git a/pkg/webservice/handlers_test.go b/pkg/webservice/handlers_test.go
index 2b2b6bdf..88553973 100644
--- a/pkg/webservice/handlers_test.go
+++ b/pkg/webservice/handlers_test.go
@@ -1207,34 +1207,21 @@ func TestGetLoggerLevel(t *testing.T) {
handler := http.HandlerFunc(getLogLevel)
handler.ServeHTTP(rr, req)
- expected := "debug"
+ expected := "info"
assert.Equal(t, rr.Body.String(), expected,
fmt.Sprintf("handler returned unexpected body: got %v want %v",
rr.Body.String(), expected))
assert.Equal(t, rr.Code, http.StatusOK)
}
func TestSetLoggerLevel(t *testing.T) {
- // invalid
req, err := http.NewRequest("PUT", "/ws/v1/loglevel",
strings.NewReader(""))
assert.NilError(t, err)
req = req.WithContext(context.WithValue(req.Context(),
httprouter.ParamsKey, httprouter.Params{
- httprouter.Param{Key: "level", Value: "invalid"},
+ httprouter.Param{Key: "level", Value: "error"},
}))
rr := httptest.NewRecorder()
-
handler := http.HandlerFunc(setLogLevel)
handler.ServeHTTP(rr, req)
- assert.Equal(t, rr.Code, http.StatusBadRequest)
-
- // valid
- req, err = http.NewRequest("PUT", "/ws/v1/loglevel",
strings.NewReader(""))
- assert.NilError(t, err)
-
- req = req.WithContext(context.WithValue(req.Context(),
httprouter.ParamsKey, httprouter.Params{
- httprouter.Param{Key: "level", Value: "error"},
- }))
- rr = httptest.NewRecorder()
- handler.ServeHTTP(rr, req)
assert.Equal(t, rr.Code, http.StatusOK)
}
diff --git a/pkg/webservice/state_dump.go b/pkg/webservice/state_dump.go
index 6d9a7dda..1c605f0f 100644
--- a/pkg/webservice/state_dump.go
+++ b/pkg/webservice/state_dump.go
@@ -67,7 +67,7 @@ func doStateDump(w io.Writer) error {
partitionContext := schedulerContext.GetPartitionMapClone()
records := imHistory.GetRecords()
- zapConfig := yunikornLog.GetConfig()
+ zapConfig := yunikornLog.GetZapConfigs()
var aggregated = AggregatedStateInfo{
Timestamp: time.Now().UnixNano(),
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]