Update internal/autoscaling/cluster_autoscaler.go
This commit is contained in:
1 parent
17bc83d221
commit
891fcd04c1
1 file changed
+108
-23
@@ -12,6 +12,7 @@ package autoscaling
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"math"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
@@ -20,21 +21,97 @@ import (
|
||||
"futriis/internal/log"
|
||||
)
|
||||
|
||||
// =============================================================================
|
||||
// СОВМЕСТИМЫЕ АТОМАРНЫЕ ОБЁРТКИ (Go 1.13+)
|
||||
// =============================================================================
|
||||
// В Go 1.19 появились atomic.Bool, atomic.Int32, atomic.Int64,
|
||||
// atomic.Uint32, atomic.Uint64, atomic.Float64. Для совместимости
|
||||
// с более старыми версиями Go (в частности, на OpenIndiana, где
|
||||
// пакет Go может быть устаревшим) реализуем собственные обёртки.
|
||||
|
||||
// atomicBool — совместимая замена atomic.Bool.
|
||||
type atomicBool struct {
|
||||
v int32
|
||||
}
|
||||
|
||||
func (b *atomicBool) Load() bool { return atomic.LoadInt32(&b.v) != 0 }
|
||||
func (b *atomicBool) Store(val bool) {
|
||||
atomic.StoreInt32(&b.v, boolToInt32(val))
|
||||
}
|
||||
func (b *atomicBool) CompareAndSwap(old, new bool) bool {
|
||||
return atomic.CompareAndSwapInt32(&b.v, boolToInt32(old), boolToInt32(new))
|
||||
}
|
||||
|
||||
func boolToInt32(b bool) int32 {
|
||||
if b {
|
||||
return 1
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
// atomicInt32 — совместимая замена atomic.Int32.
|
||||
type atomicInt32 struct {
|
||||
v int32
|
||||
}
|
||||
|
||||
func (a *atomicInt32) Load() int32 { return atomic.LoadInt32(&a.v) }
|
||||
func (a *atomicInt32) Store(val int32) { atomic.StoreInt32(&a.v, val) }
|
||||
func (a *atomicInt32) Add(delta int32) int32 { return atomic.AddInt32(&a.v, delta) }
|
||||
func (a *atomicInt32) CompareAndSwap(old, new int32) bool {
|
||||
return atomic.CompareAndSwapInt32(&a.v, old, new)
|
||||
}
|
||||
|
||||
// atomicInt64 — совместимая замена atomic.Int64.
|
||||
type atomicInt64 struct {
|
||||
v int64
|
||||
}
|
||||
|
||||
func (a *atomicInt64) Load() int64 { return atomic.LoadInt64(&a.v) }
|
||||
func (a *atomicInt64) Store(val int64) { atomic.StoreInt64(&a.v, val) }
|
||||
func (a *atomicInt64) Add(delta int64) int64 { return atomic.AddInt64(&a.v, delta) }
|
||||
|
||||
// atomicUint64 — совместимая замена atomic.Uint64.
|
||||
type atomicUint64 struct {
|
||||
v uint64
|
||||
}
|
||||
|
||||
func (a *atomicUint64) Load() uint64 { return atomic.LoadUint64(&a.v) }
|
||||
func (a *atomicUint64) Store(val uint64) { atomic.StoreUint64(&a.v, val) }
|
||||
func (a *atomicUint64) Add(delta uint64) uint64 { return atomic.AddUint64(&a.v, delta) }
|
||||
|
||||
// atomicFloat64 — совместимая замена atomic.Float64.
|
||||
// Реализована через atomic.Uint64 и math.Float64bits.
|
||||
type atomicFloat64 struct {
|
||||
v uint64
|
||||
}
|
||||
|
||||
func (f *atomicFloat64) Load() float64 {
|
||||
return math.Float64frombits(atomic.LoadUint64(&f.v))
|
||||
}
|
||||
func (f *atomicFloat64) Store(val float64) {
|
||||
atomic.StoreUint64(&f.v, math.Float64bits(val))
|
||||
}
|
||||
|
||||
// =============================================================================
|
||||
// ClusterAutoscaler
|
||||
// =============================================================================
|
||||
|
||||
// ClusterAutoscaler управляет автоматическим масштабированием кластера
|
||||
type ClusterAutoscaler struct {
|
||||
config *config.AutoscalingConfig
|
||||
logger *log.Logger
|
||||
mu sync.RWMutex
|
||||
stopChan chan struct{}
|
||||
stopOnce sync.Once
|
||||
wg sync.WaitGroup
|
||||
running atomic.Bool
|
||||
running atomicBool
|
||||
stats *AutoscalingStats
|
||||
|
||||
// Метрики нагрузки
|
||||
cpuLoad atomic.Float64
|
||||
memLoad atomic.Float64
|
||||
qps atomic.Uint64
|
||||
connections atomic.Uint64
|
||||
cpuLoad atomicFloat64
|
||||
memLoad atomicFloat64
|
||||
qps atomicUint64
|
||||
connections atomicUint64
|
||||
|
||||
// Состояние кластера
|
||||
currentNodeCount int
|
||||
@@ -48,11 +125,11 @@ type ClusterAutoscaler struct {
|
||||
|
||||
// AutoscalingStats статистика автомасштабирования
|
||||
type AutoscalingStats struct {
|
||||
TotalScaleUps atomic.Uint64
|
||||
TotalScaleDowns atomic.Uint64
|
||||
LastScaleUpAt atomic.Int64
|
||||
LastScaleDownAt atomic.Int64
|
||||
LastEvaluationAt atomic.Int64
|
||||
TotalScaleUps atomicUint64
|
||||
TotalScaleDowns atomicUint64
|
||||
LastScaleUpAt atomicInt64
|
||||
LastScaleDownAt atomicInt64
|
||||
LastEvaluationAt atomicInt64
|
||||
CurrentNodes int
|
||||
TargetNodes int
|
||||
mu sync.RWMutex
|
||||
@@ -111,13 +188,16 @@ func (a *ClusterAutoscaler) Start() {
|
||||
}
|
||||
}
|
||||
|
||||
// Stop останавливает автомасштабирование
|
||||
// Stop останавливает автомасштабирование.
|
||||
// Безопасен для многократного вызова.
|
||||
func (a *ClusterAutoscaler) Stop() {
|
||||
if !a.running.Load() {
|
||||
return
|
||||
}
|
||||
|
||||
a.stopOnce.Do(func() {
|
||||
close(a.stopChan)
|
||||
})
|
||||
a.wg.Wait()
|
||||
a.running.Store(false)
|
||||
|
||||
@@ -128,9 +208,11 @@ func (a *ClusterAutoscaler) Stop() {
|
||||
|
||||
// ReloadConfig обновляет конфигурацию автомасштабирования
|
||||
func (a *ClusterAutoscaler) ReloadConfig(cfg *config.AutoscalingConfig) {
|
||||
a.mu.Lock()
|
||||
defer a.mu.Unlock()
|
||||
if cfg == nil {
|
||||
return
|
||||
}
|
||||
|
||||
a.mu.Lock()
|
||||
oldEnabled := a.config.Enabled
|
||||
a.config = cfg
|
||||
|
||||
@@ -139,14 +221,6 @@ func (a *ClusterAutoscaler) ReloadConfig(cfg *config.AutoscalingConfig) {
|
||||
cfg.Enabled, cfg.MinNodes, cfg.MaxNodes))
|
||||
}
|
||||
|
||||
if oldEnabled != cfg.Enabled {
|
||||
if cfg.Enabled {
|
||||
a.Start()
|
||||
} else {
|
||||
a.Stop()
|
||||
}
|
||||
}
|
||||
|
||||
if a.currentNodeCount < cfg.MinNodes {
|
||||
a.currentNodeCount = cfg.MinNodes
|
||||
a.targetNodeCount = cfg.MinNodes
|
||||
@@ -155,6 +229,15 @@ func (a *ClusterAutoscaler) ReloadConfig(cfg *config.AutoscalingConfig) {
|
||||
a.currentNodeCount = cfg.MaxNodes
|
||||
a.targetNodeCount = cfg.MaxNodes
|
||||
}
|
||||
a.mu.Unlock()
|
||||
|
||||
if oldEnabled != cfg.Enabled {
|
||||
if cfg.Enabled {
|
||||
a.Start()
|
||||
} else {
|
||||
a.Stop()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// UpdateMetrics обновляет метрики нагрузки
|
||||
@@ -323,7 +406,8 @@ func (a *ClusterAutoscaler) evaluateAndScale() {
|
||||
}
|
||||
}
|
||||
|
||||
// scaleUp выполняет масштабирование вверх
|
||||
// scaleUp выполняет масштабирование вверх.
|
||||
// Вызывается под удержанием a.mu.
|
||||
func (a *ClusterAutoscaler) scaleUp(reason string) {
|
||||
count := a.config.MaxScaleUpNodes
|
||||
if a.currentNodeCount+count > a.config.MaxNodes {
|
||||
@@ -360,7 +444,8 @@ func (a *ClusterAutoscaler) scaleUp(reason string) {
|
||||
a.stats.mu.Unlock()
|
||||
}
|
||||
|
||||
// scaleDown выполняет масштабирование вниз
|
||||
// scaleDown выполняет масштабирование вниз.
|
||||
// Вызывается под удержанием a.mu.
|
||||
func (a *ClusterAutoscaler) scaleDown(reason string) {
|
||||
count := a.config.MaxScaleDownNodes
|
||||
if a.currentNodeCount-count < a.config.MinNodes {
|
||||
|
||||
Reference in new issue
Block a user