【问题标题】:Dividing tasks among Goroutines concurrently在 Goroutine 之间同时划分任务
【发布时间】:2021-04-20 11:31:54
【问题描述】:

这段代码的作用

代码从 Postgresql 数据库中获取数据。在所有数据中,只有两个字段(会话和文本)被添加到 Task Struct

我的数据库中只有 2 个(每个)提到的数据,这意味着执行 len(task) 将返回我 2 作为输出。

现在问题出在哪里:

我创建了一个长度等于任务结构的buffered channel ch(在本例中为2)。

我指定允许的最大工作线程数,这里是 20。

下面的代码所做的是,当我将任务发送到通道时,会发送Task struct(此处为 2)中的所有元素,并且 Task 结构中的示例代码将全部打印两次(= Task 结构的长度)。示例见文末。

这个程序需要做什么

例如频道len(task) = 100有100条数据。我想将这 100 个数据分成 20 个 Goroutines,每个 Goroutines 处理 5 个数据(我不知道这是否可行,如果无效,请提供其他解决方案)。

因此,这 100 个数据将提供给 20 个工作人员,他们每个人将接收 5 个数据并与他们一起运行任务,最后通道将关闭,仅此而已。

当数据库变大时,这将很有帮助。

20 个 Worker 各自执行任务,还是让 Worker 的数量等于通道中的数据数量,哪个更好?

var wg sync.WaitGroup

type Task struct {
    FetchedSession string
    FetchedText    string
}

func FetchAllData() {

    var task []Task

    //Fetch Session from DB
    var sess []database.UserSession
    database.DB.Find(&sess)
    //Fetch CommentText from DB
    var cmt []database.CommentReq
    database.DB.Find(&cmt)

    if len(sess) == len(cmt) {
        for i := range sess {
            task = append(task, Task{FetchedSession: sess[i].Session, FetchedText: cmt[i].CommentText})
        }
    }

    //making the Task Channel
    ch := make(chan []Task, len(task))

    MAX_WORKERS := 20

    wg.Add(MAX_WORKERS)

    for i := 0; i < MAX_WORKERS; i++ {
        go func() {
            for {
                t, ok := <-ch
                if !ok {
                    wg.Done()
                    return
                }
                DoTasks(t)
            }
        }()
    }

    for i := 0; i < len(task); i++ {
        ch <- task
    }

    close(ch)
    wg.Wait()
}

//Since Total number of data in Database is 2 (rows)
//Currently this function takes all data from the channel and runs Twice
func DoTasks(t []Task) {

    //Total tasks (data) = 100
    //If Max Workers = 20, then this function will run 5 times
    //Each Goroutine will get 4 tasks from the channel
    // Get the FetchedSession and FetchedTask and do tasks

    fmt.Println(t) // This prints all data twice

    //Finish one task and continue with the second
}

例子:

Example Data:
Task{FetchedSession: "EncodedString",  FetchedText: "Hello"}
Task{FetchedSession: "ExampleString",  FetchedText: "Hi"}
//Output
EncodedString
Hello
ExampleString
Hi
EncodedString
Hello
ExampleString
Hi

【问题讨论】:

  • 工人的最佳数量取决于任务和可用资源。此外,没有理由让每个工作人员完成相同数量的任务
  • @HymnsForDisco,谢谢!那么我应该怎么做这个程序呢?
  • 您的代码遍历task,并且对于每个元素,将整个task 切片发送到要处理的通道上。大概它应该只发送每个元素一次。您还应该在 goroutine 的通道上使用 range
  • @Adrian,感谢您的帮助。您能否对我的程序进行必要的更改并在下面的答案中发布。我对 Golang 很陌生,因此对 Goroutines 和通道有很大的问题。
  • 您需要做足够多的数学运算才能使您的问题有意义。 20名工人,4个任务,每人80个任务。其他 20 项任务应该去哪里?

标签: go concurrency


【解决方案1】:
  • 更改任务通道类型。
ch := make(chan Task, len(task))

这意味着通道上传递的每个值都代表一个单个任务。

  • 简化您的频道迭代
    for i := 0; i < MAX_WORKERS; i++ {
        go func() {
            defer wg.Done()
            for t := range ch {
                DoTask(t)
            }
        }()
    }

wg.Done() 现在将在 worker 退出时运行。 range ch 将在关闭通道并消耗所有任务后停止。

  • 更改“执行”功能以匹配
func DoTask(t Task) {

关于如何选择工人数量:

Run some benchmarks 用于您的FetchAllData 函数,并尝试更改MAX_WORKERS(或将其作为参数传递)。最佳值将取决于任务以及运行该功能时的可用资源,这意味着您今天机器上的最佳价值可能不是其他人机器或明天您机器上的最佳价值。基准应该可以帮助您找到一个合适的近似范围来输入值。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-11-24
    • 1970-01-01
    • 1970-01-01
    • 2020-08-09
    • 2018-03-12
    • 1970-01-01
    相关资源
    最近更新 更多