Sobes.tech
Senior

// დავწეროთ ასინქრონული დავალებების დამუშავებელი ბიბლიოთეკა // კლიენტი გადასცემს ობიექტს (Task) მონაცემებით დავალების შესრულებისთვის, // ჩვენს მაგალითში გამოვიყენებთ ცარიელ სტრუქტურას. // დამუშავებელი ერთდროულად შეუძლია არა მეტ N დავალებას, // და არა მეტ X დავალებას შეიძლება იყოს რიგში დამუშავებისთვის. // თუ რიგში ადგილი არ არის, დაუყოვნებლივ ვუბრუნებთ შეცდომას მომხმარებელს. // დავალება მიიღება დამუშავებისთვის, თუ თავისუფალი დამუშავებლები არსებობს. // ვიმიტირებთ დამუშავების ხანგრძლივობას time.Sleep(5*time.Second) გამოყენებით. // როგორც კი ერთი დავალება დასრულდება, ვიღებთ შემდეგს რიგიდან. // თუ რიგი ცარიელია, ველით ახალი დავალებებს მომხმარებლებისგან. // ვარსკვლავით: გააფართოვეთ Task სტრუქტურა და მიიღეთ დავალებების სტატუსები ცალკე მოთხოვნით.

sobes.tech AI

პასუხი AI-სგან

go package taskhandler

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

type Task struct { ID int Status string // "queued", "processing", "done", "error" }

type Handler struct { maxWorkers int maxQueue int queue chan *Task wg sync.WaitGroup mu sync.Mutex tasks map[int]*Task workerSem chan struct{} idCounter int }

func NewHandler(maxWorkers, maxQueue int) *Handler { h := &Handler{ maxWorkers: maxWorkers, maxQueue: maxQueue, queue: make(chan *Task, maxQueue), tasks: make(map[int]*Task), workerSem: make(chan struct{}, maxWorkers), } go h.dispatch() return h }

func (h *Handler) Submit(task *Task) (int, error) { h.mu.Lock() defer h.mu.Unlock()

if len(h.queue) >= h.maxQueue {
	return 0, errors.New("queue is full")
}
h.idCounter++
task.ID = h.idCounter
task.Status = "queued"
h.tasks[task.ID] = task
h.queue <- task
return task.ID, nil

}

func (h *Handler) dispatch() { for task := range h.queue { h.workerSem <- struct{}{} // acquire worker slot h.wg.Add(1) go h.process(task) } }

func (h *Handler) process(task *Task) { defer func() { <-h.workerSem // release worker slot h.wg.Done() }()

// Обновляем статус
h.mu.Lock()
task.Status = "processing"
h.mu.Unlock()

// Имитация длительной обработки
time.Sleep(5 * time.Second)

// Завершение задачи
h.mu.Lock()
task.Status = "done"
h.mu.Unlock()

}

func (h *Handler) Status(taskID int) (string, error) { h.mu.Lock() defer h.mu.Unlock() if task, ok := h.tasks[taskID]; ok { return task.Status, nil } return "", errors.New("task not found") }

func (h *Handler) Shutdown() { close(h.queue) h.wg.Wait() }


Этот код реализует асинхронный обработчик задач с ограничением на количество одновременно обрабатываемых задач (maxWorkers) и размер очереди (maxQueue). Если очередь заполнена, Submit возвращает ошибку. Статусы задач можно получить через метод Status по ID задачи.