Files
query-orchestration/internal/serviceconfig/threadpool/config.go
T
Michael McGuinness e4bdb968ba Merged in feature/threadpoolchanges (pull request #121)
Thread Pool Changes

* commentresponses
2025-04-27 15:54:52 +00:00

112 lines
2.1 KiB
Go

package threadpool
import (
"log/slog"
"runtime"
"time"
"github.com/alitto/pond/v2"
)
type ConfigProvider interface {
GetMaxWorkers() int
GetThreadMonitorInterval() int
GetThreadPool() pond.Pool
ThreadPoolStopAndWait()
StartThreadPool(*slog.Logger)
}
type ThreadPoolConfig struct {
MaxWorkers int `env:"MAX_WORKERS"`
MonitorInterval int `env:"THREAD_MONITOR_INTERVAL"`
pool pond.Pool
stopMonitor chan struct{}
logger *slog.Logger
}
type PoolMetrics struct {
QueueUtilization float64
AverageWaitTime time.Duration
ProcessingTime time.Duration
}
func (dp *ThreadPoolConfig) StartThreadPool(logger *slog.Logger) {
dp.logger = logger
if dp.MaxWorkers <= 0 {
dp.MaxWorkers = runtime.NumCPU() * 10
}
if dp.MonitorInterval <= 0 {
dp.MonitorInterval = 120
}
dp.stopMonitor = make(chan struct{})
dp.pool = pond.NewPool(dp.MaxWorkers)
// Start the monitoring heartbeat
dp.pool.Submit(dp.monitor)
}
// Log thread pool stats
func (dp *ThreadPoolConfig) monitor() {
ticker := time.NewTicker(time.Duration(dp.MonitorInterval) * time.Second)
defer ticker.Stop()
for {
select {
case <-dp.stopMonitor:
return
case <-ticker.C:
dp.logger.Debug(
"thread pool heartbeat",
"running_workers",
dp.pool.RunningWorkers(),
"queue_size",
dp.pool.QueueSize(),
"submitted_tasks",
dp.pool.SubmittedTasks(),
"waiting_tasks",
dp.pool.WaitingTasks(),
"dropped_tasks",
dp.pool.DroppedTasks(),
)
}
}
}
func (dp *ThreadPoolConfig) GetMaxWorkers() int {
return dp.MaxWorkers
}
func (dp *ThreadPoolConfig) GetThreadMonitorInterval() int {
return dp.MonitorInterval
}
func (dp *ThreadPoolConfig) GetThreadPool() pond.Pool {
if dp.pool == nil {
dp.StartThreadPool(slog.Default())
}
return dp.pool
}
func (dp *ThreadPoolConfig) ThreadPoolStopAndWait() {
waitCh := make(chan struct{})
go func() {
close(dp.stopMonitor)
dp.pool.StopAndWait()
close(waitCh)
}()
timeout := 3 * time.Second
select {
case <-waitCh:
dp.logger.Info("thread pool shutdown complete")
case <-time.After(timeout):
dp.logger.Error("thread pool shotdown timeout")
}
}