为高度并发应用程序实现全局计数器的最佳方法?

Boc*_*jim 40 concurrency go

为高度并发的应用程序实现全局计数器的最佳方法是什么?在我的情况下,我可能有10K-20K的例程执行"工作",我想计算例程正在集体工作的项目的数量和类型...

"经典"同步编码风格如下所示:

var work_counter int

func GoWorkerRoutine() {
    for {
        // do work
        atomic.AddInt32(&work_counter,1)
    }    
}
Run Code Online (Sandbox Code Playgroud)

现在这变得更加复杂,因为我想跟踪正在完成的工作的"类型",所以我真的需要这样的东西:

var work_counter map[string]int
var work_mux sync.Mutex

func GoWorkerRoutine() {
    for {
        // do work
        work_mux.Lock()
        work_counter["type1"]++
        work_mux.Unlock()
    }    
}
Run Code Online (Sandbox Code Playgroud)

似乎应该使用渠道或类似的东西"go"优化方式:

var work_counter int
var work_chan chan int // make() called somewhere else (buffered)

// started somewher else
func GoCounterRoutine() {
    for {
        select {
            case c := <- work_chan:
                work_counter += c
                break
        }
    }
}

func GoWorkerRoutine() {
    for {
        // do work
        work_chan <- 1
    }    
}
Run Code Online (Sandbox Code Playgroud)

最后一个例子仍然缺少地图,但这很容易添加.这种风格是否会提供比简单的原子增量更好的性能?当我们谈论对全局值的并发访问与可能阻止I/O完成的事情时,我无法判断这是否或多或少复杂......

我们很感激.

2013年5月28日更新:

我测试了几个实现,结果不是我的预期,这是我的计数器源代码:

package helpers

import (
)

type CounterIncrementStruct struct {
    bucket string
    value int
}

type CounterQueryStruct struct {
    bucket string
    channel chan int
}

var counter map[string]int
var counterIncrementChan chan CounterIncrementStruct
var counterQueryChan chan CounterQueryStruct
var counterListChan chan chan map[string]int

func CounterInitialize() {
    counter = make(map[string]int)
    counterIncrementChan = make(chan CounterIncrementStruct,0)
    counterQueryChan = make(chan CounterQueryStruct,100)
    counterListChan = make(chan chan map[string]int,100)
    go goCounterWriter()
}

func goCounterWriter() {
    for {
        select {
            case ci := <- counterIncrementChan:
                if len(ci.bucket)==0 { return }
                counter[ci.bucket]+=ci.value
                break
            case cq := <- counterQueryChan:
                val,found:=counter[cq.bucket]
                if found {
                    cq.channel <- val
                } else {
                    cq.channel <- -1    
                }
                break
            case cl := <- counterListChan:
                nm := make(map[string]int)
                for k, v := range counter {
                    nm[k] = v
                }
                cl <- nm
                break
        }
    }
}

func CounterIncrement(bucket string, counter int) {
    if len(bucket)==0 || counter==0 { return }
    counterIncrementChan <- CounterIncrementStruct{bucket,counter}
}

func CounterQuery(bucket string) int {
    if len(bucket)==0 { return -1 }
    reply := make(chan int)
    counterQueryChan <- CounterQueryStruct{bucket,reply}
    return <- reply
}

func CounterList() map[string]int {
    reply := make(chan map[string]int)
    counterListChan <- reply
    return <- reply
}
Run Code Online (Sandbox Code Playgroud)

它使用通道进行写入和读取,这似乎是合乎逻辑的.

这是我的测试用例:

func bcRoutine(b *testing.B,e chan bool) {
    for i := 0; i < b.N; i++ {
        CounterIncrement("abc123",5)
        CounterIncrement("def456",5)
        CounterIncrement("ghi789",5)
        CounterIncrement("abc123",5)
        CounterIncrement("def456",5)
        CounterIncrement("ghi789",5)
    }
    e<-true
}

func BenchmarkChannels(b *testing.B) {
    b.StopTimer()
    CounterInitialize()
    e:=make(chan bool)
    b.StartTimer()

    go bcRoutine(b,e)
    go bcRoutine(b,e)
    go bcRoutine(b,e)
    go bcRoutine(b,e)
    go bcRoutine(b,e)

    <-e
    <-e
    <-e
    <-e
    <-e

}

var mux sync.Mutex
var m map[string]int
func bmIncrement(bucket string, value int) {
    mux.Lock()
    m[bucket]+=value
    mux.Unlock()
}

func bmRoutine(b *testing.B,e chan bool) {
    for i := 0; i < b.N; i++ {
        bmIncrement("abc123",5)
        bmIncrement("def456",5)
        bmIncrement("ghi789",5)
        bmIncrement("abc123",5)
        bmIncrement("def456",5)
        bmIncrement("ghi789",5)
    }
    e<-true
}

func BenchmarkMutex(b *testing.B) {
    b.StopTimer()
    m=make(map[string]int)
    e:=make(chan bool)
    b.StartTimer()

    for i := 0; i < b.N; i++ {
        bmIncrement("abc123",5)
        bmIncrement("def456",5)
        bmIncrement("ghi789",5)
        bmIncrement("abc123",5)
        bmIncrement("def456",5)
        bmIncrement("ghi789",5)
    }

    go bmRoutine(b,e)
    go bmRoutine(b,e)
    go bmRoutine(b,e)
    go bmRoutine(b,e)
    go bmRoutine(b,e)

    <-e
    <-e
    <-e
    <-e
    <-e

}
Run Code Online (Sandbox Code Playgroud)

我实现了一个简单的基准测试,在地图周围只有一个互斥体(只是测试写入),并使用并行运行的5个goroutine进行基准测试.结果如下:

$ go test --bench=. helpers
PASS
BenchmarkChannels         100000             15560 ns/op
BenchmarkMutex   1000000              2669 ns/op
ok      helpers 4.452s
Run Code Online (Sandbox Code Playgroud)

我不会期望互斥体会那么快......

进一步的想法

Nic*_*ood 19

不要在链接页面中使用sync/atomic

Package atomic提供了用于实现同步算法的低级原子内存原语.这些功能需要非常小心才能正确使用.除了特殊的低级应用程序,最好使用通道或同步包的功能来实现同步

上次我不得不这样做时,我用一个互斥量来看看你的第二个例子,看起来像你的第三个带有频道的例子.当事情变得非常繁忙时,频道代码赢了,但要确保你的频道缓冲区很大.

  • 我发现`sync/atomic`文档的引用令人不满意.当然对于一个简单的计数器,`sync/atomic`是正确的工具,而且一个通道是矫枉过正的,不是吗? (3认同)
  • @Flimzy 是的,在本地基准测试之后,“sync/atomic”比基于多通道的实现快 2.5-6 倍,使用 GOMAXPROCS 通道消费者时速度快 2.5 倍,最后的按需聚合速度比单个通道快 6 倍 (3认同)

Aeq*_*tas 17

如果你正在尝试同步一个工作池(例如允许n goroutines来处理一些工作量),那么通道是一个非常好的方法,但如果你真正需要的只是一个计数器(例如页面浏览量) )然后他们是矫枉过正的.在同步和同步/原子包在那里帮助.

import "sync/atomic"

type count32 int32

func (c *count32) increment() int32 {
    return atomic.AddInt32((*int32)(c), 1)
}

func (c *count32) get() int32 {
    return atomic.LoadInt32((*int32)(c))
}
Run Code Online (Sandbox Code Playgroud)

去游乐场示例


Dij*_*tra 8

不要害怕使用互斥锁和锁只是因为你认为它们"不适合Go".在你的第二个例子中,它绝对清楚发生了什么,这很重要.您将不得不亲自尝试看看互斥体是多么满足,以及添加复杂性是否会提高性能.

如果你确实需要提高性能,也许分片是最好的方法:http: //play.golang.org/p/uLirjskGeN

缺点是您的计数只会与分片决定时一样最新.调用可能还会有很多性能time.Since(),但是,一如既往地先测量它:)

  • 互斥体是较低级别的原语,有自己有用的利基,但它们的组合方式与 CSP 设计的通道使用方式不同。因此,如果有疑问,请使用渠道。 (2认同)

Rik*_*ing 7

使用同步/原子的另一个答案适用于页面计数器之类的东西,但不适用于向外部 API 提交唯一标识符。为此,您需要一个“增量并返回”操作,该操作只能作为 CAS 循环实现。

这是一个围绕 int32 的 CAS 循环来生成唯一的消息 ID:

import "sync/atomic"

type UniqueID struct {
    counter int32
}

func (c *UniqueID) Get() int32 {
    for {
        val := atomic.LoadInt32(&c.counter)
        if atomic.CompareAndSwapInt32(&c.counter, val, val+1) {
            return val
        }
    }
}
Run Code Online (Sandbox Code Playgroud)

要使用它,只需执行以下操作:

requestID := client.msgID.Get()
form.Set("id", requestID)
Run Code Online (Sandbox Code Playgroud)

这比通道有一个优势,因为它不需要那么多额外的空闲资源 - 使用现有的 goroutines 来请求 ID,而不是为程序需要的每个计数器使用一个 goroutine。

TODO:针对渠道进行基准测试。我将猜测通道在无争用情况下更糟,在高争用情况下更好,因为它们需要排队,而此代码只是为了赢得比赛而旋转。


Ste*_*ini 6

老问题,但我只是偶然发现了这一点,它可能会有所帮助:https : //github.com/uber-go/atomic

基本上 Uber 的工程师已经在sync/atomic包之上构建了一些不错的 util 函数

我还没有在生产中对此进行测试,但代码库非常小,大多数功能的实现都非常标准

绝对优于使用通道或基本互斥锁