前往小程序,Get更优阅读体验!
立即前往
首页
学习
活动
专区
工具
TVP
发布
社区首页 >专栏 >Semaphore信号量详解

Semaphore信号量详解

原创
作者头像
大发明家
发布2021-12-15 15:26:22
1K0
发布2021-12-15 15:26:22
举报
文章被收录于专栏:技术博客文章

数据结构

semaphore.Weighted

结构体

代码语言:txt
复制
type waiter struct {
代码语言:txt
复制
    n     int64
代码语言:txt
复制
    ready chan<- struct{} // Closed when semaphore acquired.
代码语言:txt
复制
}
代码语言:txt
复制
// NewWeighted creates a new weighted semaphore with the given
代码语言:txt
复制
// maximum combined weight for concurrent access.
代码语言:txt
复制
func NewWeighted(n int64) *Weighted {
代码语言:txt
复制
    w := &Weighted{size: n}
代码语言:txt
复制
    return w
代码语言:txt
复制
}
代码语言:txt
复制
// Weighted provides a way to bound concurrent access to a resource.
代码语言:txt
复制
// The callers can request access with a given weight.
代码语言:txt
复制
type Weighted struct {
代码语言:txt
复制
    size    int64
代码语言:txt
复制
    cur     int64
代码语言:txt
复制
    mu      sync.Mutex
代码语言:txt
复制
    waiters list.List
代码语言:txt
复制
}

一个 watier 就表示一个请求,其中n表示这次请求的资源数量(权重)。

使用 NewWeighted() 函数创建一个并发访问的最大资源数,这里 n 表示资源个数。

Weighted 字段说明

  • size 表示最大资源数量,取走时会减少,释放时会增加
  • cur 计数器,记录当前已使用资源数,值范围[0 - size]
  • mu
  • waiters 当前处于等待休眠的请求者goroutine,每个请求者请求的资源数量可能不一样,只有在请求时,可用资源数量不足时请求者才会进入请求链表,每个请求者表示一个goroutine

计数器 cur 会随着资源的获取和释放而变化,那么为什么要引入数量(权重)这个概念呢?

方法列表

方法

  • NewWighted 方法用来创建一类资源,参数 n 资源表示最大可用资源总个数;
  • Acquire 获取指定个数的资源,如果当前没有空闲资源可用,当前请求者goroutine将陷入休眠状态;
  • Release 释放资源
  • TryAcquireAcquire 一样,但当无空闲资源将直接返回false,而不阻塞。

获取 Acquire 和 TryAcquire

对于获取资源有两种方法,分别为

Acquire()

TryAcquire(),两者的区别我们上面已介绍过。

在获取和释放资源前必须先加全局锁

获取资源时根据空闲资源情况,可分为三种:

  • 有空闲资源可用,将返回nil,表示成功
  • 请求资源数量超出了初始化时指定的总数量,这个肯定永远也不可能执行成功的,所以直接返回 ctx.Err()
  • 当前空闲资源数量不足,需要等待其它goroutine对资源进行释放才可以运行,这时将当前请求者goroutine放入等待队列。 这里再根据情况而定,具体见 select 判断
代码语言:txt
复制
// Acquire acquires the semaphore with a weight of n, blocking until resources
代码语言:txt
复制
// are available or ctx is done. On success, returns nil. On failure, returns
代码语言:txt
复制
// ctx.Err() and leaves the semaphore unchanged.
代码语言:txt
复制
//
代码语言:txt
复制
// If ctx is already done, Acquire may still succeed without blocking.
代码语言:txt
复制
func (s *Weighted) Acquire(ctx context.Context, n int64) error {
代码语言:txt
复制
    // 有可用资源,直接成功返回nil
代码语言:txt
复制
    s.mu.Lock()
代码语言:txt
复制
    if s.size-s.cur >= n && s.waiters.Len() == 0 {
代码语言:txt
复制
        s.cur += n
代码语言:txt
复制
        s.mu.Unlock()
代码语言:txt
复制
        return nil
代码语言:txt
复制
    }
代码语言:txt
复制
    // 请求资源权重远远超出了设置的最大权重和,失败返回 ctx.Err()
代码语言:txt
复制
    if n > s.size {
代码语言:txt
复制
        // Don't make other Acquire calls block on one that's doomed to fail.
代码语言:txt
复制
        s.mu.Unlock()
代码语言:txt
复制
        <-ctx.Done()
代码语言:txt
复制
        return ctx.Err()
代码语言:txt
复制
    }
代码语言:txt
复制
    // 有部分资源可用,将请求者放在等待队列(头部),并通过select 实现通知其它waiters
代码语言:txt
复制
    ready := make(chan struct{})
代码语言:txt
复制
    w := waiter{n: n, ready: ready}
代码语言:txt
复制
    // 放入链表尾部,并返回放入的元素
代码语言:txt
复制
    elem := s.waiters.PushBack(w)
代码语言:txt
复制
    s.mu.Unlock()
代码语言:txt
复制
    select {
代码语言:txt
复制
    case <-ctx.Done():
代码语言:txt
复制
        // 收到外面的控制信号
代码语言:txt
复制
        err := ctx.Err()
代码语言:txt
复制
        s.mu.Lock()
代码语言:txt
复制
        select {
代码语言:txt
复制
        case <-ready:
代码语言:txt
复制
            // Acquired the semaphore after we were canceled.  Rather than trying to
代码语言:txt
复制
            // fix up the queue, just pretend we didn't notice the cancelation.
代码语言:txt
复制
            // 如果在用户取消之前已经获取了资源,则直接忽略这个信号,返回nil表示成功
代码语言:txt
复制
            err = nil
代码语言:txt
复制
        default:
代码语言:txt
复制
            // 收到控制信息,且还没有获取到资源,就直接将原来添加的 waiter 删除
代码语言:txt
复制
            isFront := s.waiters.Front() == elem
代码语言:txt
复制
            // 则将其从链接删除 上面 ctx.Done()
代码语言:txt
复制
            s.waiters.Remove(elem)
代码语言:txt
复制
            // 如果当前元素正好位于链表最前面,且还存在可用的资源,就通知其它waiters
代码语言:txt
复制
            if isFront && s.size > s.cur {
代码语言:txt
复制
                s.notifyWaiters()
代码语言:txt
复制
            }
代码语言:txt
复制
        }
代码语言:txt
复制
        s.mu.Unlock()
代码语言:txt
复制
        return err
代码语言:txt
复制
    case <-ready:
代码语言:txt
复制
        return nil
代码语言:txt
复制
    }
代码语言:txt
复制
}

注意上面在select逻辑语句上面有一次加解锁的操作,在 select 里面由于是全局锁所以还需要再次加锁。

根据可用计数器信息,可分三种情况:

  1. 对于 TryAcquire() 就比较简单了,就是一个可用资源数量的判断,数量够用表示成功返回 true ,否则 false,此方法并不会进行阻塞,而是直接返回。
代码语言:txt
复制
// TryAcquire acquires the semaphore with a weight of n without blocking.
代码语言:txt
复制
// On success, returns true. On failure, returns false and leaves the semaphore unchanged.
代码语言:txt
复制
func (s *Weighted) TryAcquire(n int64) bool {
代码语言:txt
复制
    s.mu.Lock()
代码语言:txt
复制
    success := s.size-s.cur >= n && s.waiters.Len() == 0
代码语言:txt
复制
    if success {
代码语言:txt
复制
        s.cur += n
代码语言:txt
复制
    }
代码语言:txt
复制
    s.mu.Unlock()
代码语言:txt
复制
    return success
代码语言:txt
复制
}

释放 Release

对于释放也很简单,就是将已使用资源数量(计数器)进行更新减少,并通知其它 waiters

代码语言:txt
复制
// Release releases the semaphore with a weight of n.
代码语言:txt
复制
func (s *Weighted) Release(n int64) {
代码语言:txt
复制
    s.mu.Lock()
代码语言:txt
复制
    s.cur -= n
代码语言:txt
复制
    if s.cur < 0 {
代码语言:txt
复制
        s.mu.Unlock()
代码语言:txt
复制
        panic("semaphore: released more than held")
代码语言:txt
复制
    }
代码语言:txt
复制
    s.notifyWaiters()
代码语言:txt
复制
    s.mu.Unlock()
代码语言:txt
复制
}

通知机制

通过 for 循环从链表头部开始头部依次遍历出链表中的所有waiter,并更新计数器 Weighted.cur,同时将其从链表中删除,直到遇到

空闲资源数量 < watier.n 为止。

代码语言:txt
复制
func (s *Weighted) notifyWaiters() {
代码语言:txt
复制
    for {
代码语言:txt
复制
        next := s.waiters.Front()
代码语言:txt
复制
        if next == nil {
代码语言:txt
复制
            break // No more waiters blocked.
代码语言:txt
复制
        }
代码语言:txt
复制
        w := next.Value.(waiter)
代码语言:txt
复制
        if s.size-s.cur < w.n {
代码语言:txt
复制
            // Not enough tokens for the next waiter.  We could keep going (to try to
代码语言:txt
复制
            // find a waiter with a smaller request), but under load that could cause
代码语言:txt
复制
            // starvation for large requests; instead, we leave all remaining waiters
代码语言:txt
复制
            // blocked.
代码语言:txt
复制
            //
代码语言:txt
复制
            // Consider a semaphore used as a read-write lock, with N tokens, N
代码语言:txt
复制
            // readers, and one writer.  Each reader can Acquire(1) to obtain a read
代码语言:txt
复制
            // lock.  The writer can Acquire(N) to obtain a write lock, excluding all
代码语言:txt
复制
            // of the readers.  If we allow the readers to jump ahead in the queue,
代码语言:txt
复制
            // the writer will starve — there is always one token available for every
代码语言:txt
复制
            // reader.
代码语言:txt
复制
            break
代码语言:txt
复制
        }
代码语言:txt
复制
        s.cur += w.n
代码语言:txt
复制
        s.waiters.Remove(next)
代码语言:txt
复制
        close(w.ready)
代码语言:txt
复制
    }
代码语言:txt
复制
}

可以看到如果一个链表里有多个等待者,其中一个等待者需要的资源(权重)比较多的时候,当前 watier

会出现长时间的阻塞(即使当前可用资源足够其它waiter执行,期间会有一些资源浪费), 直到有足够的资源可以让这个等待者执行,然后继续执行它后面的等待者。

使用示例

官方文档提供了一个基于信号量的典型的“工作池”模式,见[https://pkg.go.dev/golang.org/x/sync/semaphore#example-

package-

WorkerPool](https://links.jianshu.com/go?to=https%3A%2F%2Fpkg.go.dev%2Fgolang.org%2Fx%2Fsync%2Fsemaphore%23example-

package-WorkerPool),演示了如何通过信号量控制一定数量的 goroutine 并发工作。

这是一个通过信号量实现并发对

考拉兹猜想的示例,对1-32之间的数字进行计算,并打印32个符合结果的值。

代码语言:txt
复制
package main
代码语言:txt
复制
import (
代码语言:txt
复制
    "context"
代码语言:txt
复制
    "fmt"
代码语言:txt
复制
    "log"
代码语言:txt
复制
    "runtime"
代码语言:txt
复制
    "golang.org/x/sync/semaphore"
代码语言:txt
复制
)
代码语言:txt
复制
// Example_workerPool demonstrates how to use a semaphore to limit the number of
代码语言:txt
复制
// goroutines working on parallel tasks.
代码语言:txt
复制
//
代码语言:txt
复制
// This use of a semaphore mimics a typical “worker pool” pattern, but without
代码语言:txt
复制
// the need to explicitly shut down idle workers when the work is done.
代码语言:txt
复制
func main() {
代码语言:txt
复制
    ctx := context.TODO()
代码语言:txt
复制
    // 权重值为逻辑cpu个数
代码语言:txt
复制
    var (
代码语言:txt
复制
        maxWorkers = runtime.GOMAXPROCS(0)
代码语言:txt
复制
        sem        = semaphore.NewWeighted(int64(maxWorkers))
代码语言:txt
复制
        out        = make([]int, 32)
代码语言:txt
复制
    )
代码语言:txt
复制
    // Compute the output using up to maxWorkers goroutines at a time.
代码语言:txt
复制
    for i := range out {
代码语言:txt
复制
        // When maxWorkers goroutines are in flight, Acquire blocks until one of the
代码语言:txt
复制
        // workers finishes.
代码语言:txt
复制
        if err := sem.Acquire(ctx, 1); err != nil {
代码语言:txt
复制
            log.Printf("Failed to acquire semaphore: %v", err)
代码语言:txt
复制
            break
代码语言:txt
复制
        }
代码语言:txt
复制
        go func(i int) {
代码语言:txt
复制
            defer sem.Release(1)
代码语言:txt
复制
            out[i] = collatzSteps(i + 1)
代码语言:txt
复制
        }(i)
代码语言:txt
复制
    }
代码语言:txt
复制
    // 如果使用了 errgroup 原语则不需要下面这段语句
代码语言:txt
复制
    if err := sem.Acquire(ctx, int64(maxWorkers)); err != nil {
代码语言:txt
复制
        log.Printf("Failed to acquire semaphore: %v", err)
代码语言:txt
复制
    }
代码语言:txt
复制
    fmt.Println(out)
代码语言:txt
复制
}
代码语言:txt
复制
// collatzSteps computes the number of steps to reach 1 under the Collatz
代码语言:txt
复制
// conjecture. (See https://en.wikipedia.org/wiki/Collatz_conjecture.)
代码语言:txt
复制
func collatzSteps(n int) (steps int) {
代码语言:txt
复制
    if n <= 0 {
代码语言:txt
复制
        panic("nonpositive input")
代码语言:txt
复制
    }
代码语言:txt
复制
    for ; n > 1; steps++ {
代码语言:txt
复制
        if steps < 0 {
代码语言:txt
复制
            panic("too many steps")
代码语言:txt
复制
        }
代码语言:txt
复制
        if n%2 == 0 {
代码语言:txt
复制
            n /= 2
代码语言:txt
复制
            continue
代码语言:txt
复制
        }
代码语言:txt
复制
        const maxInt = int(^uint(0) >> 1)
代码语言:txt
复制
        if n > (maxInt-1)/3 {
代码语言:txt
复制
            panic("overflow")
代码语言:txt
复制
        }
代码语言:txt
复制
        n = 3*n + 1
代码语言:txt
复制
    }
代码语言:txt
复制
    return steps
代码语言:txt
复制
}

上面先声明了总权重值为逻辑CPU数量,每次 for 循环都会调用一次 sem.Acquire(ctx, 1), 即表示最多每个CPU可运行一个

goroutine,如果当前权重值不足的话,其它groutine将处于阻塞状态,这里共循环32次,即阻塞数量最大为 32-maxWorkers

每获取成功一个权重就会执行go匿名函数,并在函数结束时释放权重。为了保证每次for循环都会正常结束,最后调用了 `sem.Acquire(ctx,

int64(maxWorkers)),表示最后一次执行必须获取的权重值为maxWorkers。当然如果使用errgroup`

同步原语的话,这一步可以省略掉

以下为使用 errgroup 的方法

代码语言:txt
复制
func main() {
代码语言:txt
复制
    ctx := context.TODO()
代码语言:txt
复制
    var (
代码语言:txt
复制
        maxWorkers = runtime.GOMAXPROCS(0)
代码语言:txt
复制
        sem        = semaphore.NewWeighted(int64(maxWorkers))
代码语言:txt
复制
        out        = make([]int, 32)
代码语言:txt
复制
    )
代码语言:txt
复制
    group, _ := errgroup.WithContext(context.Background())
代码语言:txt
复制
    for i := range out {
代码语言:txt
复制
        if err := sem.Acquire(ctx, 1); err != nil {
代码语言:txt
复制
            log.Printf("Failed to acquire semaphore: %v", err)
代码语言:txt
复制
            break
代码语言:txt
复制
        }
代码语言:txt
复制
        group.Go(func() error {
代码语言:txt
复制
            go func(i int) {
代码语言:txt
复制
                defer sem.Release(1)
代码语言:txt
复制
                out[i] = collatzSteps(i + 1)
代码语言:txt
复制
            }(i)
代码语言:txt
复制
            return nil
代码语言:txt
复制
        })
代码语言:txt
复制
    }
代码语言:txt
复制
    // 这里会阻塞,直到所有goroutine都执行完毕
代码语言:txt
复制
    if err := group.Wait(); err != nil {
代码语言:txt
复制
        fmt.Println(err)
代码语言:txt
复制
    }
代码语言:txt
复制
    fmt.Println(out)
代码语言:txt
复制
}

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

评论
作者已关闭评论
0 条评论
热度
最新
推荐阅读
目录
  • 数据结构
  • 方法列表
  • 获取 Acquire 和 TryAcquire
  • 释放 Release
  • 通知机制
  • 使用示例
领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档