120
}
121
123
f.logger.Error("Target worker pool size is negative. Please fix the dynamic config.", tag.Key("worker-pool-size"), tag.Value(targetWorkerNum))
124
return
125
}
126
128
>
if targetWorkerNum == currentWorkerNum {
129
return
130
}
131
133
>
f.startWorkers(targetWorkerNum - currentWorkerNum)
134
>
} else {
135
f.stopWorkers(currentWorkerNum - targetWorkerNum)
136
}
137
138
>
f.logger.Info("Update worker pool size", tag.Key("worker-pool-size"), tag.Value(targetWorkerNum))
fifo_scheduler.go
139
}
140
141
func (f *FIFOScheduler[T]) startWorkers(
142
count int,
144
>
for range count {
145
>
shutdownCh := make(chan struct{})
146
>
f.workerShutdownCh = append(f.workerShutdownCh, shutdownCh)
147
>
148
>
f.shutdownWG.Add(1)
149
>
go f.processTask(shutdownCh)
150
>
}
151
}
152