752fb2e2c0
Text Extraction Clean Up + Parallelization * merged_Call * directchildren * bitofimprovements * go * parallel * fixfullsuite * pondforconcurrency * comments * muchdone * bitsofclean * stabilisedtests * snappy * threadpooltests * childelements * testspassed
112 lines
2.1 KiB
Go
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()
|
|
SetThreadPool(*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) SetThreadPool(logger *slog.Logger) {
|
|
dp.logger = logger
|
|
|
|
if dp.MaxWorkers <= 0 {
|
|
dp.MaxWorkers = runtime.NumCPU() * 10000
|
|
}
|
|
|
|
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.SetThreadPool(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")
|
|
}
|
|
}
|