【问题标题】:What is the Advantage of sync.WaitGroup over Channels?与 Channels 相比,sync.WaitGroup 的优势是什么?
【发布时间】:2016-07-03 13:39:04
【问题描述】:

我正在开发一个并发 Go 库,我偶然发现了两种不同的 goroutine 之间的同步模式,它们的结果相似:

Waitgroup

package main

import (
    "fmt"
    "sync"
    "time"
)

var wg sync.WaitGroup

func main() {
    words := []string{"foo", "bar", "baz"}

    for _, word := range words {
        wg.Add(1)
        go func(word string) {
            time.Sleep(1 * time.Second)
            defer wg.Done()
            fmt.Println(word)
        }(word)
    }
    // do concurrent things here

    // blocks/waits for waitgroup
    wg.Wait()
}

Channel

package main

import (
    "fmt"
    "time"
)

func main() {
    words := []string{"foo", "bar", "baz"}
    done := make(chan bool)
    // defer close(done)
    for _, word := range words {
        // fmt.Println(len(done), cap(done))
        go func(word string) {
            time.Sleep(1 * time.Second)
            fmt.Println(word)
            done <- true
        }(word)
    }
    // Do concurrent things here

    // This blocks and waits for signal from channel
    for range words {
        <-done
    }
}

我被告知sync.WaitGroup 的性能稍好一些,而且我已经看到它被普遍使用。但是,我发现频道更惯用。与频道相比,使用sync.WaitGroup 的真正优势是什么和/或更好的情况可能是什么情况?

【问题讨论】:

  • 在您的第二个示例中,同步错误。你阻塞直到第一个 goroutine 在通道上发送,而不是直到最后一个。
  • 真正地道,大多数“爆炸”通道(仅用于发送信号的通道)应该具有chan struct{} 类型而不是chan bool。此外,频道在下面使用sync,因此使用sync 应该更高效。 WaitGroup 在您必须阻止等待许多 goroutine 返回时提供帮助。如果您可以在 for 循环中生成 100 个它们,那就更简单了。
  • 我尝试运行基于频道的代码,但没有成功。更正的版本在这里。 play.golang.org/p/LHx8Tto-kvI。使用等待组是惯用的,但是如果您想控制并发性,我会担心如何使用等待组来做到这一点。是否有限制等待组。使用通道你可以做到这一点。有一个缓冲通道,然后在该过程完成后读取通道。这样就可以处理下一个项目。
  • @Angelo,我已经更正了你的代码:play.golang.org/p/CglhQg0eVjL(三个 goroutine 没有同时运行,并且始终按此顺序打印“foo bar baz”。)

标签: go concurrency channel


【解决方案1】:

与您的第二个示例的正确性无关(如 cmets 中所述,您并没有按照您的想法做,但它很容易修复),我倾向于认为第一个示例更容易掌握。

现在,我什至不会说频道更惯用。通道是 Go 语言的一个标志性功能,并不意味着尽可能地使用它们是惯用的。 Go 中的惯用语是使用最简单和最容易理解的解决方案:在这里,WaitGroup 传达了含义(您的主要功能是 Waiting 让工人完成)和机械师(工人在他们完成时通知是Done)。

除非您处于非常特殊的情况,否则我不建议在此处使用渠道解决方案。

【讨论】:

    【解决方案2】:

    这取决于用例。如果您正在调度一次性作业以并行运行而无需知道每个作业的结果,那么您可以使用WaitGroup。但是如果你需要从 goroutines 中收集结果,那么你应该使用一个通道。

    由于通道是双向的,所以我几乎总是使用通道。

    另一方面,正如评论中指出的那样,您的频道示例未正确实施。您将需要一个单独的频道来指示没有更多的工作要做(一个例子是here)。在您的情况下,由于您事先知道单词的数量,因此您可以只使用一个缓冲通道并接收固定次数以避免声明关闭通道。

    【讨论】:

      【解决方案3】:

      对于您的简单示例(表示作业完成),WaitGroup 是显而易见的选择。 Go 编译器非常友好,不会责怪您使用通道来完成任务的简单信号,但某些代码审查员会这样做。

      1. “WaitGroup 等待一组 goroutine 完成。 主 goroutine 调用Add(n) 设置数量 等待的 goroutines。然后每个 goroutine 完成后运行并调用Done()。同时, Wait 可用于阻塞,直到所有 goroutine 完成。”
      words := []string{"foo", "bar", "baz"}
      var wg sync.WaitGroup
      for _, word := range words {
          wg.Add(1)
          go func(word string) {
              defer wg.Done()
              time.Sleep(100 * time.Millisecond) // a job
              fmt.Println(word)
          }(word)
      }
      wg.Wait()
      

      可能性仅限于您的想象力:

      1. 频道可以缓冲
      words := []string{"foo", "bar", "baz"}
      done := make(chan struct{}, len(words))
      for _, word := range words {
          go func(word string) {
              time.Sleep(100 * time.Millisecond) // a job
              fmt.Println(word)
              done <- struct{}{} // not blocking
          }(word)
      }
      for range words {
          <-done
      }
      
      1. 通道可以无缓冲,您可以只使用一个信号通道(例如chan struct{}):
      words := []string{"foo", "bar", "baz"}
      done := make(chan struct{})
      for _, word := range words {
          go func(word string) {
              time.Sleep(100 * time.Millisecond) // a job
              fmt.Println(word)
              done <- struct{}{} // blocking
          }(word)
      }
      for range words {
          <-done
      }
      
      1. 您可以限制具有缓冲通道容量的并发作业数:
      t0 := time.Now()
      var wg sync.WaitGroup
      words := []string{"foo", "bar", "baz"}
      done := make(chan struct{}, 1) // set the number of concurrent job here
      for _, word := range words {
          wg.Add(1)
          go func(word string) {
              done <- struct{}{}
              time.Sleep(100 * time.Millisecond) // job
              fmt.Println(word, time.Since(t0))
              <-done
              wg.Done()
          }(word)
      }
      wg.Wait()
      
      1. 您可以使用频道发送消息:
      done := make(chan string)
      go func() {
          for _, word := range []string{"foo", "bar", "baz"} {
              done <- word
          }
          close(done)
      }()
      for word := range done {
          fmt.Println(word)
      }
      

      基准测试:

          go test -benchmem -bench . -args -n 0
      # BenchmarkEvenWaitgroup-8  1827517   652 ns/op    0 B/op  0 allocs/op
      # BenchmarkEvenChannel-8    1000000  2373 ns/op  520 B/op  1 allocs/op
          go test -benchmem -bench .
      # BenchmarkEvenWaitgroup-8  1770260   678 ns/op    0 B/op  0 allocs/op
      # BenchmarkEvenChannel-8    1560124  1249 ns/op  158 B/op  0 allocs/op
      

      代码(main_test.go):

      package main
      
      import (
          "flag"
          "fmt"
          "os"
          "sync"
          "testing"
      )
      
      func BenchmarkEvenWaitgroup(b *testing.B) {
          evenWaitgroup(b.N)
      }
      func BenchmarkEvenChannel(b *testing.B) {
          evenChannel(b.N)
      }
      func evenWaitgroup(n int) {
          if n%2 == 1 { // make it even:
              n++
          }
          for i := 0; i < n; i++ {
              wg.Add(1)
              go func(n int) {
                  select {
                  case ch <- n: // tx if channel is empty
                  case i := <-ch: // rx if channel is not empty
                      // fmt.Println(n, i)
                      _ = i
                  }
                  wg.Done()
              }(i)
          }
          wg.Wait()
      }
      func evenChannel(n int) {
          if n%2 == 1 { // make it even:
              n++
          }
          for i := 0; i < n; i++ {
              go func(n int) {
                  select {
                  case ch <- n: // tx if channel is empty
                  case i := <-ch: // rx if channel is not empty
                      // fmt.Println(n, i)
                      _ = i
                  }
                  done <- struct{}{}
              }(i)
          }
          for i := 0; i < n; i++ {
              <-done
          }
      }
      func TestMain(m *testing.M) {
          var n int // We use TestMain to set up the done channel.
          flag.IntVar(&n, "n", 1_000_000, "chan cap")
          flag.Parse()
          done = make(chan struct{}, n)
          fmt.Println("n=", n)
          os.Exit(m.Run())
      }
      
      var (
          done chan struct{}
          ch   = make(chan int)
          wg   sync.WaitGroup
      )
      

      【讨论】:

      • 是的。我认为使用 Channels 替换 Waitgroups 是反模式。 Channels 和 WaitGroups 是不可替换的,但它们被设计为相互集成。通过这种方式,您的解决方案非常出色。
      【解决方案4】:

      如果您特别坚持只使用渠道,那么它需要以不同的方式进行(如果我们使用您的示例,正​​如@Not_a_Golfer 指出的那样,它会产生不正确的结果)。

      一种方法是创建一个 int 类型的通道。在工作进程中,每次完成作业时发送一个数字(这也可以是唯一的作业 ID,如果您愿意,可以在接收器中跟踪它)。

      在接收者主 go 例程中(它将知道提交的确切作业数) - 在通道上执行范围循环,直到提交的作业数未完成,并在所有作业时退出循环完成。如果您想跟踪每个作业的完成情况(如果需要,可能会做一些事情),这是一个好方法。

      这是供您参考的代码。减少 totalJobsLeft 将是安全的,因为它只会在通道的范围循环中完成!

      //This is just an illustration of how to sync completion of multiple jobs using a channel
      //A better way many a times might be to use wait groups
      
      package main
      
      import (
          "fmt"
          "math/rand"
          "time"
      )
      
      func main() {
      
          comChannel := make(chan int)
          words := []string{"foo", "bar", "baz"}
      
          totalJobsLeft := len(words)
      
          //We know how many jobs are being sent
      
          for j, word := range words {
              jobId := j + 1
              go func(word string, jobId int) {
      
                  fmt.Println("Job ID:", jobId, "Word:", word)
                  //Do some work here, maybe call functions that you need
                  //For emulating this - Sleep for a random time upto 5 seconds
                  randInt := rand.Intn(5)
                  //fmt.Println("Got random number", randInt)
                  time.Sleep(time.Duration(randInt) * time.Second)
                  comChannel <- jobId
              }(word, jobId)
          }
      
          for j := range comChannel {
              fmt.Println("Got job ID", j)
              totalJobsLeft--
              fmt.Println("Total jobs left", totalJobsLeft)
              if totalJobsLeft == 0 {
                  break
              }
          }
          fmt.Println("Closing communication channel. All jobs completed!")
          close(comChannel)
      
      }
      

      【讨论】:

        【解决方案5】:

        我经常使用通道从可能产生错误的 goroutine 中收集错误消息。这是一个简单的例子:

        func couldGoWrong() (err error) {
            errorChannel := make(chan error, 3)
        
            // start a go routine
            go func() (err error) {
                defer func() { errorChannel <- err }()
        
                for c := 0; c < 10; c++ {
                    _, err = fmt.Println(c)
                    if err != nil {
                        return
                    }
                }
        
                return
            }()
        
            // start another go routine
            go func() (err error) {
                defer func() { errorChannel <- err }()
        
                for c := 10; c < 100; c++ {
                    _, err = fmt.Println(c)
                    if err != nil {
                        return
                    }
                }
        
                return
            }()
        
            // start yet another go routine
            go func() (err error) {
                defer func() { errorChannel <- err }()
        
                for c := 100; c < 1000; c++ {
                    _, err = fmt.Println(c)
                    if err != nil {
                        return
                    }
                }
        
                return
            }()
        
            // synchronize go routines and collect errors here
            for c := 0; c < cap(errorChannel); c++ {
                err = <-errorChannel
                if err != nil {
                    return
                }
            }
        
            return
        }
        

        【讨论】:

          【解决方案6】:

          也建议使用waitgroup,但你还是想用channel来做,那么下面我提到一个简单的channel使用

          package main
          
          import (
              "fmt"
              "time"
          )
          
          func main() {
              c := make(chan string)
              words := []string{"foo", "bar", "baz"}
          
              go printWordrs(words, c)
          
              for j := range c {
                  fmt.Println(j)
              }
          }
          
          
          func printWordrs(words []string, c chan string) {
              defer close(c)
              for _, word := range words {
                  time.Sleep(1 * time.Second)
                  c <- word
              }   
          }
          

          【讨论】:

          • 你在哪里关闭了频道?
          猜你喜欢
          • 2016-07-04
          • 2011-03-13
          • 1970-01-01
          • 1970-01-01
          • 2011-09-27
          • 1970-01-01
          • 2013-06-16
          • 2010-10-20
          • 1970-01-01
          相关资源
          最近更新 更多