【问题标题】:Golang: newbie -- Master-Worker concurrencyGolang:新手——Master-Worker并发
【发布时间】:2013-12-09 01:20:52
【问题描述】:

我在尝试实现这一点时遇到了问题(所有 goroutines 都睡着了 - 死锁!) 这是代码的要点:

var workers = runtime.NumCPU()

func main() {
    jobs := make(chan *myStruct, workers)
    done := make(chan *myStruct, workers)

    go produceWork(file_with_jobs, jobs)
    for i := 0; i < runtime.NumCPU(); i++ {
        go Worker(jobs, done)
    }
    consumeWork(done)
}

func produceWork(vf string, jobs chan *utils.DigSigEntries) {
    defer close(jobs)

    // load file with jobs
    file, err := ini.LoadFile(vf)

    // get data for processing
    for data, _ := range file {
        // ...
        jobs <- &myStruct{data1, data2, data3, false}
    }
}

func Worker(in, out chan *myStruct) {
    for {
        item, open := <-in
        if !open {
            break
        }

        process(item)
        out <- item
    }
    // close(out)   --> tried closing the out channel, but then not all items are processed
    //                  though no panics occur.
}

func process(item *myStruct) {
    //...modify the item
    item.status = true
}

func consumeWork(done chan *myStruct) {
    for val := range done {
        if !val.status {
            fmt.Println(val)
        }
    }
}

我主要是想了解如何在不使用同步/等待的东西的情况下做到这一点 - 只是纯频道 - 这可能吗?此例程的目标是让单个生产者加载由 N 个工人处理的项目 - 感谢任何指针/帮助。

【问题讨论】:

  • 如果其他 goroutine 将在通道上写入或读取,您无法关闭通道。一旦关闭,它就对所有人关闭。您需要另一种机制来跟踪您的活动 goroutine,例如,sync.WaitGroup,通过通道发送特殊值,使用另一个通道作为控制通道(通过它发送“我快死了”)...

标签: go


【解决方案1】:

您可以按照 siritinga 的建议,使用第三个信号或计数器通道,例如signal chan boolean,其中produceWork goroutine 会在每个作业进入jobs 通道之前添加一个值。因此,与jobs 一样,将向signal 传递相同数量的值:

func produceWork(vf string, jobs chan *utils.DigSigEntries, signal chan boolean) {
    defer close(jobs)

    // load file with jobs
    file, err := ini.LoadFile(vf)

    // get data for processing
    for data, _ := range file {
        // ...
        signal <- true
        jobs <- &myStruct{data1, data2, data3, false}
    }

    close(signal)
}

消费将从signal 频道开始读取。如果有一个值,则可以肯定会有一个从out 通道读取的值(一旦工作人员将其传递)。如果signal 关闭,那么一切都完成了。我们可以关闭剩余的done频道:

func consumeWork(done chan *myStruct, signal chan boolean) {
    for _ := range signal {
        val <- done
        if !val.status {
            fmt.Println(val)
        }
    }

    close(done)
} 

虽然这是可能的,但我不会真的推荐它。它不会使代码比使用sync.WaitGroup 时更清晰。毕竟,signal 频道基本上只能用作计数器。 WaitGroup 具有相同的目的,而且成本更低。

但您的问题不是关于如何解决问题,而是是否可以通过纯渠道解决问题。

【讨论】:

    【解决方案2】:

    抱歉,我没有注意到您想跳过 /sync :/ 我会留下答案,也许有人正在寻找这个。

    import (
        "sync"
    )
    
    
    func main() {
        jobs := make(chan *myStruct, workers)
        done := make(chan *myStruct, workers)
    
        var workerWg sync.WaitGroup   // waitGroup for workers
        var consumeWg sync.WaitGroup  // waitGroup for consumer
        consumeWg.Add(1) // add one active Consumer
    
        for i := 0; i < runtime.NumCPU(); i++ {
            go Worker(&workerWg, jobs, done)
            workerWg.Add(1)
        }
        go consumeWork(&consumeWg, done)
    
        produceWork(file_with_jobs, jobs)
    
        close(jobs)
        workerWg.Wait()
    
        close(done)
        consumeWg.Wait()
    
    }
    
    func produceWork(vf string, jobs chan *utils.DigSigEntries) {
    
        // load file with jobs
        file, err := ini.LoadFile(vf)
    
        // get data for processing
        for data, _ := range file {
            // ...
            jobs <- &myStruct{data1, data2, data3, false}
        }
    }
    
    
    func Worker(wg *sync.WaitGroup, done chan *myStruct) {
    
        defer wg.Done()
        for job := range jobs {
    
            result := process(job)
            out <- result
        }
    
       // close(out)   --> tried closing the out channel, but then not all items are processed
       //                  though no panics occur.
    }
    
    func process(item *myStruct) {
        //...modify the item
        item.status = true
    }
    
    func consumeWork(wg *sync.WaitGroup, done chan *myStruct) {
        defer wg.Done()
        for val := range done {
            if !val.status {
                fmt.Println(val)
            }
        }
    }
    

    【讨论】:

    • 我同意在这种情况下使用同步可能是最好的。但是您可以只用一个wgJobs sync.WaitGroup 来简化它。在将每个作业添加到通道之前让produceWork Add() (consumeWork 将在每个作业被消耗后调用 Done())。当produceWork 没有更多工作要添加时,close(jobs); wgJobs.Wait(); close(done);
    猜你喜欢
    • 2016-07-28
    • 1970-01-01
    • 2018-11-17
    • 1970-01-01
    • 1970-01-01
    • 2016-05-24
    • 2016-05-11
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多