first
This commit is contained in:
@@ -0,0 +1,115 @@
|
||||
package control_case
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io/fs"
|
||||
"log"
|
||||
"os"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/gocarina/gocsv"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
func exists(path string) (bool, error) {
|
||||
_, err := os.Stat(path)
|
||||
if err == nil {
|
||||
return true, nil
|
||||
}
|
||||
if errors.Is(err, fs.ErrNotExist) {
|
||||
return false, nil
|
||||
}
|
||||
return false, err
|
||||
}
|
||||
|
||||
func worker(id int, jobs <-chan Params, ctx context.Context) {
|
||||
for j := range jobs {
|
||||
fmt.Printf("Worker %d started job %d \n", id, j.ID)
|
||||
j.CreateFolder()
|
||||
j.Run(ctx)
|
||||
time.Sleep(time.Second) // Имитация длительной задачи
|
||||
fmt.Printf("Worker %d finished job %d \n", id, j.ID)
|
||||
// results <- j * 2
|
||||
}
|
||||
}
|
||||
|
||||
func readCSVWithGocsv() []Params {
|
||||
file, err := os.Open("db.csv")
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
defer file.Close()
|
||||
|
||||
var params []Params
|
||||
if err := gocsv.UnmarshalFile(file, ¶ms); err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
|
||||
for _, person := range params {
|
||||
fmt.Printf("%+v\n", person)
|
||||
}
|
||||
return params
|
||||
}
|
||||
|
||||
func readFromDb(db *gorm.DB) []ControlCase {
|
||||
var cases []ControlCase
|
||||
statuses := []string{"N"}
|
||||
_ = db.
|
||||
Where("status IN ?", statuses).
|
||||
Preload("Params").
|
||||
Find(&cases).Error
|
||||
return cases
|
||||
}
|
||||
|
||||
func RunControllWorker(db *gorm.DB, results chan int) {
|
||||
const numJobs = 1
|
||||
const numWorkers = 10
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer func() {
|
||||
fmt.Printf("cencel")
|
||||
cancel()
|
||||
}() // Cancel if main exits early
|
||||
|
||||
jobs := make(chan Params, numJobs)
|
||||
//results := make(chan int, numJobs)
|
||||
|
||||
// Запуск воркеров
|
||||
var wg sync.WaitGroup
|
||||
for w := 1; w <= numWorkers; w++ {
|
||||
wg.Add(1)
|
||||
go func(w int) {
|
||||
defer wg.Done()
|
||||
worker(w, jobs, ctx)
|
||||
}(w)
|
||||
}
|
||||
|
||||
//params := readCSVWithGocsv()
|
||||
for {
|
||||
cases := readFromDb(db)
|
||||
// Отправка задач
|
||||
for _, c := range cases {
|
||||
p := c.Params
|
||||
e, _ := p.Exist()
|
||||
if e {
|
||||
c.Status = "R"
|
||||
db.Save((&c))
|
||||
} else {
|
||||
p.CreateFolder()
|
||||
jobs <- p
|
||||
}
|
||||
}
|
||||
time.Sleep(10 * time.Second)
|
||||
}
|
||||
// Ожидание завершения всех воркеров
|
||||
//fmt.Printf("wait")
|
||||
// wg.Wait()
|
||||
// close(jobs)
|
||||
//
|
||||
// // Получение результатов
|
||||
// for r := range results {
|
||||
// fmt.Println("Result:", r)
|
||||
// }
|
||||
}
|
||||
Reference in New Issue
Block a user