Обработка []int через worker pool для CPU-bound задачи
Есть CPU-bound функция heavyCompute(n int) int, обрабатывающая один int. Реализуйте processAll(nums []int) []int, прогоняющую через неё весь срез параллельно через worker pool. Порядок результата сохраняется — out[i] соответствует nums[i]. Требования: размер пула — по числу ядер (runtime.NumCPU()), а не по одной goroutine на элемент; раздавайте работу через канал; синхронизируйтесь через sync.WaitGroup.
func processAll(nums []int) []int {
// ваш код здесь
}
Допишите реализацию.
Для CPU-bound работы размер пула — runtime.NumCPU(): лишние goroutine сверх числа ядер дают лишь накладные расходы планировщика. Раздавайте индексы через канал задач; каждый воркер пишет out[i] = heavyCompute(nums[i]) в свой индекс — блокировка не нужна, порядок сохраняется. sync.WaitGroup ждёт завершения всех воркеров.
- ✗Запускать по goroutine на элемент для CPU-bound работы — тысячи goroutine забивают планировщик вместо насыщения ядер
- ✗Собирать результаты через
appendв общий срез под mutex — это сериализует воркеров и путает порядок вывода - ✗Брать размер пула по числу элементов или хардкод-константой вместо
runtime.NumCPU()
- →Почему воркеров больше, чем ядер CPU, не ускоряет CPU-bound работу?
- →Как заменить этот ручной пул на пакет
errgroupсSetLimit?
Скелет для реализации
CPU-bound работу нет смысла дробить мельче числа ядер: больше goroutine, чем ядер, не ускоряют счёт, а лишь нагружают планировщик. Раздача индексов через канал задач динамически балансирует нагрузку, если стоимость элементов неравномерна. Каждый воркер пишет в свой индекс out[i] — поэтому ни блокировка, ни общий append не нужны, и порядок сохраняется сам собой.
func processAll(nums []int) []int {
out := make([]int, len(nums))
jobs := make(chan int, len(nums)) // индексы элементов
var wg sync.WaitGroup
workers := runtime.NumCPU() // по числу ядер, не по элементу
for w := 0; w < workers; w++ {
wg.Add(1)
go func() {
defer wg.Done()
for i := range jobs { // забирает следующий свободный индекс
out[i] = heavyCompute(nums[i]) // пишет в свой индекс — без блокировок
}
}()
}
for i := range nums {
jobs <- i
}
close(jobs) // воркеры выйдут из range после слива канала
wg.Wait() // ждём завершения всех
return out
}