Sobes.tech
Senior

// Müştəri, yerinə yetirmək üçün məlumatlar ilə birlikdə (Task) obyektini ötürür, // nümunəmizdə boş bir struktur istifadə edəcəyik. // // Planlayıcı maksimum N vəzifəni eyni zamanda işləyə bilər, // və X vəzifə işləmə üçün növbədə ola bilər. // Əgər növbədə yer yoxdursa, dərhal müştəriyə xəta qaytarılır. // // Vəzifə, sərbəst emalçılar mövcuddursa, qəbul edilir. // İşləmə müddətini time.Sleep(5*time.Second) ilə simulyasiya edirik. // Hər vəzifə tamamlandıqda, növbəti vəzifəni növbədən alırıq. // Əgər növbə boşdursa, yeni vəzifələr üçün gözləyirik. // // Ulduz işarəsi ilə: Task strukturunu genişləndirin və vəzifələrin statuslarını ayrıca sorğu ilə alın. package scheduler type Scheduler interface { AddTask(t Task) error Close() } type scheduler struct { queue chan Task wg sync.WaitGroup } func New(workers, queue int) Scheduler { // gözlənilir implementasiya } type Task interface { Do() }

sobes.tech Süni İntellekt

AI-dan cavab

Ниже пример реализации Scheduler в Go, который ограничивает количество одновременно обрабатываемых задач (workers) и размер очереди (queue). Если очередь заполнена, возвращается ошибка. Задачи обрабатываются параллельно, после завершения задачи берётся следующая из очереди.

package scheduler

import (
	"errors"
	"sync"
	"time"
)

type Task interface {
	Do()
}

type scheduler struct {
	queue   chan Task
	workers int
	wg      sync.WaitGroup
	quit    chan struct{}
}

func New(workers, queue int) Scheduler {
	s := &scheduler{
		queue:   make(chan Task, queue),
		workers: workers,
		quit:    make(chan struct{}),
	}
	for i := 0; i < workers; i++ {
		go s.worker()
	}
	return s
}

func (s *scheduler) worker() {
	for {
		select {
		case task := <-s.queue:
			s.wg.Add(1)
			task.Do()
			s.wg.Done()
		case <-s.quit:
			return
		}
	}
}

func (s *scheduler) AddTask(t Task) error {
	select {
	case s.queue <- t:
		return nil
	default:
		return errors.New("queue is full")
	}
}

func (s *scheduler) Close() {
	close(s.quit)
	// Ждём завершения всех задач
	s.wg.Wait()
}

// Пример реализации Task
// type MyTask struct{}
// func (t MyTask) Do() {
//     time.Sleep(5 * time.Second) // имитация работы
// }

Пояснения:

  • queue — буферизированный канал для задач, размером queue.
  • Если канал заполнен, AddTask возвращает ошибку.
  • workers — количество горутин-обработчиков.
  • Каждая горутина читает задачи из канала и выполняет Do().
  • Close останавливает обработчики и ждёт завершения всех задач.

Для реализации статусов задач можно расширить Task структурой с ID и статусом, а также добавить мапу для хранения статусов с синхронизацией.