【问题标题】:Why not all my goroutines executed? Need an explaination [closed]为什么我的所有 goroutine 不执行?需要解释[关闭]
【发布时间】:2021-03-03 02:24:27
【问题描述】:
    func main() {
    var wg sync.WaitGroup
    ch := make(chan int,5)
    start:= time.Now()
    cnt:=0
    wg.Add(10000)
    for i:=0; i<10000; i++ {
        ch <- 1
        go func() {
            defer wg.Done()
            doSomething()
            cnt++
            <-ch
        }()
    }
    wg.Wait()
    fmt.Println(cnt)
    end:= time.Now()
    fmt.Println("End of program.",end.Sub(start))
}

这里我想同时执行程序,并且我希望最多有 5 个 goroutine。 问题是当我打印出“cnt”时,它不会是 10000。这意味着我有一些未执行的 goroutine。我该如何解决这个问题?

现在我正在使用互斥锁来解决问题,但是这个程序的运行时间不会因为 goroutine 而更好,我不明白为什么。

    func main() {
    var wg sync.WaitGroup
    var mutex sync.Mutex
    ch := make(chan int,5)
    start:= time.Now()
    cnt:=0
    for i:=0; i<10000; i++ {
        wg.Add(1)
        ch <- 1
        go func() {
            defer wg.Done()
            defer mutex.Unlock()
            mutex.Lock()
            doSomething()
            cnt++
            <-ch
        }()
    }
    wg.Wait()
    fmt.Println(cnt)
    end:= time.Now()
    fmt.Println("End of program.",end.Sub(start))
}

【问题讨论】:

  • 你有一个数据竞赛。结果毫无意义。
  • 你在那里操作一个变量并且在并发期间它是不安全的,尝试在变量和wg.add(1)循环之间使用sycn.Mutex
  • 使用互斥锁保护cnt++
  • 您更新的代码现在在 Goroutine 开始时锁定一个互斥锁,并在完成后解锁它;这将阻止并发执行。如果doSomething() 是线程安全的,那么您只需要在cnt++ 周围保持锁定,即mutex.Lock(); cnt++; mutex.Unlock()

标签: go concurrency goroutine


【解决方案1】:

尝试在以下程序中使用工人的数量,看看哪个数字能给你带来最好的结果。不过,您进行基准测试的方式并不可靠,通常应该避免使用。

但是这个程序应该更好用;肯定有更好的实现。所以在这里你只是产生 workers 数量的 goroutines 但在你的情况下,它是 10000 goroutines 这对于小情况来说真的没有必要;这是一个矫枉过正。

注意:对我来说,这个程序比你的实现好 50% 以上。

package main

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

type work struct {
    wg   *sync.WaitGroup
    jobs <-chan struct{}

    mu  *sync.Mutex
    val *int
}

// worker is responsible for executing the work assigned
func worker(w work) {
    for range w.jobs {
        w.mu.Lock()
        *w.val++
        w.mu.Unlock()
    }
    w.wg.Done()
}

func main() {
    start := time.Now()
    jobs := make(chan struct{}, 2) // Number of jobs (buffer)
    workers := 2                   // Number of workers
    cnt := 0                       // Shared variable among workers

    work := work{
        wg:   &sync.WaitGroup{},
        jobs: jobs,
        mu:   &sync.Mutex{},
        val:  &cnt,
    }

    // Worker Pool
    work.wg.Add(workers)
    for i := 0; i < workers; i++ {
        go worker(work)
    }

    // Allocate jobs (Signal worker(s))
    for i := 0; i < 10000; i++ {
        jobs <- struct{}{}
    }
    // Ask the workers to stop
    close(jobs)
    work.wg.Wait()

    fmt.Println(cnt)
    fmt.Println("End of program: ", time.Since(start))
}

【讨论】:

    猜你喜欢
    • 2014-08-17
    • 1970-01-01
    • 2017-12-15
    • 2013-03-24
    • 1970-01-01
    • 2011-05-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多