mirror of
https://github.com/DullJZ/s3-balance.git
synced 2026-09-04 23:36:40 +08:00
Support hot update config
This commit is contained in:
@@ -87,7 +87,7 @@ func (b *Balancer) GetStrategy() string {
|
||||
// 允许在运行时更改负载均衡策略
|
||||
func (b *Balancer) SetStrategy(strategyName string) error {
|
||||
var strategy Strategy
|
||||
|
||||
|
||||
switch strategyName {
|
||||
case "round-robin":
|
||||
strategy = NewRoundRobinStrategy()
|
||||
@@ -100,11 +100,16 @@ func (b *Balancer) SetStrategy(strategyName string) error {
|
||||
default:
|
||||
return fmt.Errorf("unknown balancer strategy: %s", strategyName)
|
||||
}
|
||||
|
||||
|
||||
b.strategy = strategy
|
||||
return nil
|
||||
}
|
||||
|
||||
// UpdateStrategy 更新负载均衡策略(热更新用)
|
||||
func (b *Balancer) UpdateStrategy(strategyName string) error {
|
||||
return b.SetStrategy(strategyName)
|
||||
}
|
||||
|
||||
// SetMetrics 设置指标服务
|
||||
func (b *Balancer) SetMetrics(metrics *metrics.Metrics) {
|
||||
b.metrics = metrics
|
||||
|
||||
+100
-1
@@ -3,6 +3,7 @@ package bucket
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
@@ -280,7 +281,7 @@ func (m *Manager) GetVirtualBuckets() []*BucketInfo {
|
||||
func (m *Manager) GetRealBuckets() []*BucketInfo {
|
||||
m.mu.RLock()
|
||||
defer m.mu.RUnlock()
|
||||
|
||||
|
||||
var real []*BucketInfo
|
||||
for _, b := range m.buckets {
|
||||
if !b.IsVirtual() {
|
||||
@@ -289,3 +290,101 @@ func (m *Manager) GetRealBuckets() []*BucketInfo {
|
||||
}
|
||||
return real
|
||||
}
|
||||
|
||||
// UpdateConfig 更新配置(支持热更新)
|
||||
func (m *Manager) UpdateConfig(newConfig *config.Config) error {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
|
||||
log.Println("Updating bucket manager configuration...")
|
||||
|
||||
// 更新配置引用
|
||||
oldConfig := m.config
|
||||
m.config = newConfig
|
||||
|
||||
// 检查是否需要重新创建存储桶
|
||||
needsRestart := m.checkIfRestartNeeded(oldConfig, newConfig)
|
||||
if needsRestart {
|
||||
log.Println("Bucket configuration changed significantly, recreating buckets...")
|
||||
|
||||
// 停止现有的监控
|
||||
if m.healthMonitor != nil {
|
||||
m.healthMonitor.Stop()
|
||||
}
|
||||
if m.statsMonitor != nil {
|
||||
m.statsMonitor.Stop()
|
||||
}
|
||||
|
||||
// 重新创建bucket映射
|
||||
m.buckets = make(map[string]*BucketInfo)
|
||||
|
||||
// 初始化所有存储桶客户端
|
||||
for _, bucketCfg := range newConfig.Buckets {
|
||||
if !bucketCfg.Enabled {
|
||||
continue
|
||||
}
|
||||
|
||||
client, err := createS3Client(bucketCfg)
|
||||
if err != nil {
|
||||
// 如果创建失败,恢复旧配置
|
||||
m.config = oldConfig
|
||||
return fmt.Errorf("failed to create S3 client for bucket %s: %v", bucketCfg.Name, err)
|
||||
}
|
||||
|
||||
info := &BucketInfo{
|
||||
Config: bucketCfg,
|
||||
Client: client,
|
||||
Available: true,
|
||||
LastChecked: time.Now(),
|
||||
}
|
||||
|
||||
m.buckets[bucketCfg.Name] = info
|
||||
}
|
||||
|
||||
// 重新初始化监控
|
||||
m.initHealthMonitoring()
|
||||
} else {
|
||||
// 只更新监控间隔(需要重启监控来改变间隔)
|
||||
log.Println("Updating monitoring intervals...")
|
||||
if m.healthMonitor != nil {
|
||||
m.healthMonitor.Stop()
|
||||
}
|
||||
if m.statsMonitor != nil {
|
||||
m.statsMonitor.Stop()
|
||||
}
|
||||
// 重新初始化监控以应用新的间隔
|
||||
m.initHealthMonitoring()
|
||||
}
|
||||
|
||||
log.Println("Bucket manager configuration updated successfully")
|
||||
return nil
|
||||
}
|
||||
|
||||
// checkIfRestartNeeded 检查是否需要重启bucket manager
|
||||
func (m *Manager) checkIfRestartNeeded(oldConfig, newConfig *config.Config) bool {
|
||||
// 检查bucket数量变化
|
||||
if len(oldConfig.Buckets) != len(newConfig.Buckets) {
|
||||
return true
|
||||
}
|
||||
|
||||
// 检查bucket配置变化
|
||||
for i, oldBucket := range oldConfig.Buckets {
|
||||
if i >= len(newConfig.Buckets) {
|
||||
return true
|
||||
}
|
||||
newBucket := newConfig.Buckets[i]
|
||||
|
||||
// 检查关键配置项
|
||||
if oldBucket.Name != newBucket.Name ||
|
||||
oldBucket.Endpoint != newBucket.Endpoint ||
|
||||
oldBucket.AccessKeyID != newBucket.AccessKeyID ||
|
||||
oldBucket.SecretAccessKey != newBucket.SecretAccessKey ||
|
||||
oldBucket.Region != newBucket.Region ||
|
||||
oldBucket.Enabled != newBucket.Enabled ||
|
||||
oldBucket.Virtual != newBucket.Virtual {
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
return false
|
||||
}
|
||||
|
||||
@@ -0,0 +1,246 @@
|
||||
package config
|
||||
|
||||
import (
|
||||
"log"
|
||||
"os"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/fsnotify/fsnotify"
|
||||
)
|
||||
|
||||
// Manager 配置管理器,支持热更新
|
||||
type Manager struct {
|
||||
configFile string
|
||||
config *Config
|
||||
mutex sync.RWMutex
|
||||
watcher *fsnotify.Watcher
|
||||
callbacks []func(*Config)
|
||||
stopChan chan struct{}
|
||||
lastModTime time.Time
|
||||
pollingTicker *time.Ticker
|
||||
}
|
||||
|
||||
// NewManager 创建新的配置管理器
|
||||
func NewManager(configFile string) (*Manager, error) {
|
||||
// 初始加载配置
|
||||
cfg, err := Load(configFile)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// 获取文件的初始修改时间
|
||||
fileInfo, err := os.Stat(configFile)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
manager := &Manager{
|
||||
configFile: configFile,
|
||||
config: cfg,
|
||||
callbacks: make([]func(*Config), 0),
|
||||
stopChan: make(chan struct{}),
|
||||
lastModTime: fileInfo.ModTime(),
|
||||
}
|
||||
|
||||
// 同时启用fsnotify和轮询监听
|
||||
// 这样可以确保在Docker挂载等场景下也能正常工作
|
||||
manager.initWatching()
|
||||
|
||||
return manager, nil
|
||||
}
|
||||
|
||||
// initWatching 初始化文件监听(同时使用fsnotify和轮询)
|
||||
func (m *Manager) initWatching() {
|
||||
// 尝试启用fsnotify
|
||||
watcher, err := fsnotify.NewWatcher()
|
||||
if err == nil {
|
||||
if err := watcher.Add(m.configFile); err == nil {
|
||||
m.watcher = watcher
|
||||
log.Println("fsnotify watcher enabled for config file")
|
||||
go m.watchConfig()
|
||||
} else {
|
||||
log.Printf("Failed to add file to fsnotify watcher: %v", err)
|
||||
watcher.Close()
|
||||
}
|
||||
} else {
|
||||
log.Printf("Failed to create fsnotify watcher: %v", err)
|
||||
}
|
||||
|
||||
// 同时启用轮询模式(作为备用和补充)
|
||||
// 在Docker挂载等场景下,轮询更可靠
|
||||
m.pollingTicker = time.NewTicker(3 * time.Second)
|
||||
log.Println("Config file polling enabled (3s interval)")
|
||||
go m.pollConfig()
|
||||
}
|
||||
|
||||
// pollConfig 轮询检查配置文件变化
|
||||
func (m *Manager) pollConfig() {
|
||||
for {
|
||||
select {
|
||||
case <-m.pollingTicker.C:
|
||||
fileInfo, err := os.Stat(m.configFile)
|
||||
if err != nil {
|
||||
log.Printf("Failed to stat config file during polling: %v", err)
|
||||
continue
|
||||
}
|
||||
|
||||
// 检查文件修改时间
|
||||
if fileInfo.ModTime().After(m.lastModTime) {
|
||||
log.Printf("Config file %s modified (detected by polling), reloading...", m.configFile)
|
||||
m.lastModTime = fileInfo.ModTime()
|
||||
m.reloadConfig()
|
||||
}
|
||||
|
||||
case <-m.stopChan:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// GetConfig 获取当前配置(线程安全)
|
||||
func (m *Manager) GetConfig() *Config {
|
||||
m.mutex.RLock()
|
||||
defer m.mutex.RUnlock()
|
||||
|
||||
// 返回配置的副本以避免并发修改
|
||||
configCopy := *m.config
|
||||
return &configCopy
|
||||
}
|
||||
|
||||
// OnConfigChange 注册配置变化回调
|
||||
func (m *Manager) OnConfigChange(callback func(*Config)) {
|
||||
m.mutex.Lock()
|
||||
defer m.mutex.Unlock()
|
||||
m.callbacks = append(m.callbacks, callback)
|
||||
}
|
||||
|
||||
// watchConfig 监听配置文件变化(fsnotify模式)
|
||||
func (m *Manager) watchConfig() {
|
||||
for {
|
||||
select {
|
||||
case event, ok := <-m.watcher.Events:
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
|
||||
// 只处理修改和重命名事件
|
||||
if event.Op&fsnotify.Write == fsnotify.Write ||
|
||||
event.Op&fsnotify.Rename == fsnotify.Rename {
|
||||
log.Printf("Config file %s modified (detected by fsnotify), reloading...", m.configFile)
|
||||
|
||||
// 更新最后修改时间以避免轮询重复触发
|
||||
if fileInfo, err := os.Stat(m.configFile); err == nil {
|
||||
m.lastModTime = fileInfo.ModTime()
|
||||
}
|
||||
|
||||
m.reloadConfig()
|
||||
}
|
||||
|
||||
case err, ok := <-m.watcher.Errors:
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
log.Printf("Config watcher error: %v", err)
|
||||
|
||||
case <-m.stopChan:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// reloadConfig 重新加载配置
|
||||
func (m *Manager) reloadConfig() {
|
||||
// 添加延迟以防止编辑器的多次写入事件
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
|
||||
// 加载新配置
|
||||
newConfig, err := Load(m.configFile)
|
||||
if err != nil {
|
||||
log.Printf("Failed to reload config: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
// 更新配置
|
||||
m.mutex.Lock()
|
||||
oldConfig := m.config
|
||||
m.config = newConfig
|
||||
callbacks := make([]func(*Config), len(m.callbacks))
|
||||
copy(callbacks, m.callbacks)
|
||||
m.mutex.Unlock()
|
||||
|
||||
log.Printf("Configuration reloaded successfully")
|
||||
|
||||
// 异步调用回调函数
|
||||
go func() {
|
||||
for _, callback := range callbacks {
|
||||
func() {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
log.Printf("Config change callback panic: %v", r)
|
||||
}
|
||||
}()
|
||||
callback(newConfig)
|
||||
}()
|
||||
}
|
||||
}()
|
||||
|
||||
// 记录重要配置变更
|
||||
m.logConfigChanges(oldConfig, newConfig)
|
||||
}
|
||||
|
||||
// logConfigChanges 记录配置变更
|
||||
func (m *Manager) logConfigChanges(oldConfig, newConfig *Config) {
|
||||
// 检查服务器端口变化
|
||||
if oldConfig.Server.Port != newConfig.Server.Port {
|
||||
log.Printf("Server port changed: %d -> %d (restart required)",
|
||||
oldConfig.Server.Port, newConfig.Server.Port)
|
||||
}
|
||||
|
||||
// 检查数据库配置变化
|
||||
if oldConfig.Database.DSN != newConfig.Database.DSN {
|
||||
log.Printf("Database DSN changed (restart required)")
|
||||
}
|
||||
|
||||
// 检查存储桶数量变化
|
||||
if len(oldConfig.Buckets) != len(newConfig.Buckets) {
|
||||
log.Printf("Bucket count changed: %d -> %d",
|
||||
len(oldConfig.Buckets), len(newConfig.Buckets))
|
||||
}
|
||||
|
||||
// 检查负载均衡策略变化
|
||||
if oldConfig.Balancer.Strategy != newConfig.Balancer.Strategy {
|
||||
log.Printf("Load balancer strategy changed: %s -> %s",
|
||||
oldConfig.Balancer.Strategy, newConfig.Balancer.Strategy)
|
||||
}
|
||||
|
||||
// 检查代理模式变化
|
||||
if oldConfig.S3API.ProxyMode != newConfig.S3API.ProxyMode {
|
||||
log.Printf("S3 API proxy mode changed: %t -> %t",
|
||||
oldConfig.S3API.ProxyMode, newConfig.S3API.ProxyMode)
|
||||
}
|
||||
|
||||
// 检查指标配置变化
|
||||
if oldConfig.Metrics.Enabled != newConfig.Metrics.Enabled {
|
||||
log.Printf("Metrics enabled changed: %t -> %t",
|
||||
oldConfig.Metrics.Enabled, newConfig.Metrics.Enabled)
|
||||
}
|
||||
}
|
||||
|
||||
// Close 关闭配置管理器
|
||||
func (m *Manager) Close() error {
|
||||
// 停止监听协程
|
||||
close(m.stopChan)
|
||||
|
||||
// 停止轮询
|
||||
if m.pollingTicker != nil {
|
||||
m.pollingTicker.Stop()
|
||||
}
|
||||
|
||||
// 关闭fsnotify watcher
|
||||
if m.watcher != nil {
|
||||
return m.watcher.Close()
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
Reference in New Issue
Block a user