Middle
Workshop چیست؟
sobes.tech AI
Ҷавоб аз AI
Работещо крадене (work stealing) е механизъм за балансиране на натоварването в многопоточни системи, при който неактивният поток (работник) "краде" задачи от опашката на друг, зает поток.
Ключови характеристики:
- Децентрализирано: Всеки поток управлява своята локална опашка с задачи.
- Активно: Неактивният поток активно търси работа при други.
- Висока производителност: Ефективен за задачи с неравномерно натоварване.
В Golang се използва в планировчика (scheduler) за разпределяне на горутините по потоци на операционната система (M). Когато един M завърши текущата си работа (пакет горутин на P), той търси горутини при други P чрез техните локални опашки.
// Псевдокод, демонстриращ концепцията за крадене на работни задачи
package main
import "fmt"
import "sync"
import "time"
type Task struct {
id int
}
// Имитация на опашка с задачи
type TaskQueue struct {
tasks []Task
mu sync.Mutex
}
func (q *TaskQueue) AddTask(task Task) {
q.mu.Lock()
defer q.mu.Unlock()
q.tasks = append(q.tasks, task)
}
func (q *TaskQueue) GetLocalTask() (Task, bool) {
q.mu.Lock()
defer q.mu.Unlock()
if len(q.tasks) == 0 {
return Task{}, false
}
task := q.tasks[0]
q.tasks = q.tasks[1:]
return task, true
}
// Опит за "крадене" на задача от друга опашка
func (q *TaskQueue) StealTask() (Task, bool) {
q.mu.Lock()
defer q.mu.Unlock()
if len(q.tasks) < 2 { // Не "крадем", ако има малко задачи
return Task{}, false
}
// Крадем от средата или края, за да избегнем конфликт с локалното извличане
index := len(q.tasks) / 2
task := q.tasks[index]
q.tasks = append(q.tasks[:index], q.tasks[index+1:]...)
return task, true
}
// Имитация на работно потокче
func worker(id int, localQueue *TaskQueue, otherQueues []*TaskQueue, wg *sync.WaitGroup) {
defer wg.Done()
for {
// Опит за вземане на задача от локалната опашка
task, ok := localQueue.GetLocalTask()
if ok {
fmt.Printf("Работник %d изпълнява задача %d локално\n", id, task.id)
time.Sleep(100 * time.Millisecond) // Имитация на работа
continue
}
// Ако локалната опашка е празна, опитваме да крадем
stolen := false
for _, queue := range otherQueues {
if queue == localQueue {
continue // Не крадем от себе си
}
task, ok := queue.StealTask()
if ok {
fmt.Printf("Работник %d краде задача %d\n", id, task.id)
time.Sleep(100 * time.Millisecond) // Имитация на работа
stolen = true
break // Успешно крадене, излизаме
}
}
if !stolen {
// Ако не е успял да краде, може да няма повече задачи
// В реален планировчик има по-сложни механизми за приключване
// За демонстрация просто излизаме след известно време
fmt.Printf("Работник %d е в режим на изчакване...\n", id)
time.Sleep(50 * time.Millisecond)
// В реален сценарий тук би имало механизъм за приключване или паркиране на потока
// return // За демонстрация позволяваме изход
break // За опростяване на симулацията
}
}
}
func main() {
numWorkers := 4
numTasks := 20
taskQueues := make([]*TaskQueue, numWorkers)
for i := range taskQueues {
taskQueues[i] = &TaskQueue{}
}
// Разпределяме задачите неравномерно за демонстрация на крадене
for i := 0; i < numTasks; i++ {
queueIndex := i % 2 // Повече задачи за първите две работни
taskQueues[queueIndex].AddTask(Task{id: i})
}
var wg sync.WaitGroup
for i := 0; i < numWorkers; i++ {
wg.Add(1)
// Предаваме референции към всички опашки, за да може работникът да "краде"
otherQueues := make([]*TaskQueue, 0, numWorkers-1)
for j := range taskQueues {
if i != j {
otherQueues = append(otherQueues, taskQueues[j])
}
}
go worker(i, taskQueues[i], otherQueues, &wg)
}
wg.Wait()
fmt.Println("Всички задачи са изпълнени")
}