channel

不要通过共享内存来通信,而要通过通信来共享内存

以下是 go 1.24.4版本的 channel 的实现,尝试去解读

数据结构

chan结构

hchan

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
type hchan struct {
qcount uint // total data in the queue
dataqsiz uint // size of the circular queue
buf unsafe.Pointer // points to an array of dataqsiz elements
elemsize uint16
synctest bool // true if created in a synctest bubble
closed uint32
timer *timer // timer feeding this chan
elemtype *_type // element type
sendx uint // send index
recvx uint // receive index
recvq waitq // list of recv waiters
sendq waitq // list of send waiters
lock mutex
}

channel 的数据结构 hchan

  1. qcount 当前 channel 缓冲区存在的数据数量
  2. dataqsiz 当前缓冲区能存放的数据数量
  3. buf 指向当前 channel 用来存放数据的缓冲区(一个环形数组)
  4. elemsize 存放的数据类型的大小
  5. synctest 用于标记该 channel 是否在 testing/synctest 的 bubble 中创建
  6. closed 用来标识当前 channel 是否关闭
  7. timer 如果这个 channel 是 time 包生成的定时器通道,那这里指向负责在到点时往它里面写入 time.Time 并配合调度器管理触发/阻塞的那个 runtime 定时器对象
  8. elemtype 存放的数据的类型
  9. sendx 发送数据到环形缓冲区的 index
  10. recvx 从环形缓冲区读取数据的 index
  11. recvq 等待向 channel 读取数据的队列
  12. sendq 等待向 channel 发送数据的队列
  13. lock 互斥锁

waitq

1
2
3
4
type waitq struct {
first *sudog
last *sudog
}

waitq 相当于封装了一个等待双向链表,维护该双向链表的头尾指针,内部用 sudog 的 next/prev 组成双向链表

sudog

1
2
3
4
5
6
7
8
9
10
11
type sudog struct {
g *g
next *sudog
prev *sudog
elem unsafe.Pointer
isSelect bool
success bool
// ...
waitlink *sudog
c *hchan
}

sudog 双向等待链表

  1. g 当前 sudog 节点对应的 goroutine
  2. next 指向下一个 sudog 节点
  3. prev 指向上一个 sudog 节点
  4. elem 数据元素
  5. isSelect 是否在 select 多路复用中
  6. success 标识是由于通道通信成功唤醒,还是被 channel 管道关闭唤醒
  7. waitlink 一个 goroutine 可能对应的多个 sudog 在多个 channel 中加入等待队列,waitlink 用于将这些 sudog 串联起来挂在 g.waiting,方便一个 case 执行后,沿着 waitlink 将每个 sudog 从各自 channel 的 next/prev 双向链表中删掉。
  8. c 与当前 sudog 交互的 chan

demo_select

1
2
3
4
5
6
7
8
9
10
go func() {
select {
case <-ch1:

case <-ch2:

case x := <-ch3:

}
}()

当三个 channel 的操作都需要阻塞时,当前 goroutine 对应的多个 sudog 会分别挂到三个 channel 的等待队列上,同时这些 sudog 通过 waitlink 串成一条链挂在 g.waiting,方便某个 case 成功后一次性从其他 channel 的等待队列中删除对应的 sudog。

创建 channel

两种方式: 带缓冲区或者不带缓冲区

1
2
ch1 := make(chan struct{}) 
ch2 := make(chan struct{}, size)

ch1是创建的不带缓冲区的 channel,ch2是创建的缓冲区大小为 size 的 channel,上述两个创建 channel 的方式在编译器阶段会被类型检查标记为“创建 channel 的内建 make”,从而变为对 runtime.makechan 或者 runtime.makechan64 的一次调用,runtime.makechan64只是封装了对 runtime.makechan 的调用,所以主要看 runtime.makechan 方法。

makechan

makechan(t * chantype, size int)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
func makechan(t *chantype, size int) *hchan {
elem := t.Elem
// ...
mem, overflow := math.MulUintptr(elem.Size_, uintptr(size))
if overflow || mem > maxAlloc-hchanSize || size < 0 {
panic(plainError("makechan: size out of range"))
}
var c *hchan
switch {
case mem == 0:
// Queue or element size is zero.
c = (*hchan)(mallocgc(hchanSize, nil, true))
// Race detector uses this location for synchronization.
c.buf = c.raceaddr()
case !elem.Pointers():
// Elements do not contain pointers.
// Allocate hchan and buf in one call.
c = (*hchan)(mallocgc(hchanSize+mem, nil, true))
c.buf = add(unsafe.Pointer(c), hchanSize)
default:
// Elements contain pointers.
c = new(hchan)
c.buf = mallocgc(mem, elem, true)
}

c.elemsize = uint16(elem.Size_)
c.elemtype = elem
c.dataqsiz = uint(size)
if getg().syncGroup != nil {
c.synctest = true
}
lockInit(&c.lock, lockRankHchan)

//...
return c
}

该方法有两个入参,分别是要创建的 channel 的类型信息和元素个数

  1. 首先取出要创建 channel 的元素类型,判断元素类型的大小和元素个数相乘获取到申请内存的大小 mem 是否越界,当无缓冲区时元素个数为0,则 mem 为0
  2. 初始化 chan 时,会根据不同类型不同处理方式,分别是,当 mem == 0 时(无缓冲或元素为 struct{} 等零大小类型);当 mem != 0 且元素不含指针(PtrBytes == 0);当 mem != 0 且元素包含指针(PtrBytes > 0)。
  3. 如果 mem == 0 时则直接申请一个大小为默认值 hchanSize 的内存空间,具体是几字节(96、104 等)取决于当前版本的 hchan 字段布局和平台对齐要求。
  4. 当 mem != 0且元素不包含指针则申请一个 hchanSize+mem 的连续内存空间,并且将 c.buf 指向缓冲区的起始地址
  5. 当 mem != 0且元素包含指针,则分别申请 chan 和 buf 的空间
  6. 接着完成对 channel 其他的字段初始化,最后返回 channel 指针

channel 的读取和发送数据的阻塞和非阻塞

在 runtime.channel.go 中有阻塞和非阻塞的两种执行逻辑,通过 bool 类型的 block 参数实现

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
//	select {
// case c <- v:
// ... foo
// default:
// ... bar
// }
//
// as
//
// if selected, ok = selectnbrecv(&v, c); selected {
// ... foo
// } else {
// ... bar
// }
func selectnbrecv(elem unsafe.Pointer, c *hchan) (selected, received bool) {
return chanrecv(c, elem, false)
}
// select {
// case v, ok = <-c:
// ... foo
// default:
// ... bar
// }
//
// as
//
// if selectnbsend(c, v) {
// ... foo
// } else {
// ... bar
// }
func selectnbsend(c *hchan, elem unsafe.Pointer) (selected bool) {
return chansend(c, elem, false, sys.GetCallerPC())
}

在默认情况,读/写 channel 都是阻塞模式,但是在 go 官方给的注解可以看到在 select 多路复用时,读/写 channel 操作会被汇编为 selectnbrecv / selectnbsend 方法,本质还是调用 chanrecv / chansend 方法,但是此时传入的 block 为 false 代表非阻塞读写,实现了非阻塞模式。

写入数据

chansend1(c * hchan, elem unsafe.Pointer)

1
2
3
func chansend1(c *hchan, elem unsafe.Pointer) {
chansend(c, elem, true, sys.GetCallerPC())
}

chansend1主要是封装了对 chansend 的调用,下面阐述 chansend 方法,由于 chansend 过长所以分段阐述。

case1

send1

chansend(c * hchan, ep unsafe.Pointer, block bool, callerpc uintptr)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
func chansend(c *hchan, ep unsafe.Pointer, block bool, callerpc uintptr) bool {
if c == nil {
if !block {
return false
}
gopark(nil, nil, waitReasonChanSendNilChan, traceBlockForever, 2)
throw("unreachable")
}
//..

if !block && c.closed == 0 && full(c) {
return false
}
//..
}

func full(c *hchan) bool {
if c.dataqsiz == 0 {
return c.recvq.first == nil
}
return c.qcount == c.dataqsiz
}

block 参数标识是否以阻塞形式发送数据,首先先看快速处理:

  1. 如果向一个未初始化的 channel 发送数据并且 block 为 true,则 gopark ,否则直接返回 false
  2. 接下来如果以非阻塞形式向未关闭的 channel 发送数据,则会调用 full 判断 channel 现在能不能再接收一个 send 而不阻塞,不可以的话直接返回 false

full 方法会根据不同类型进行不同逻辑判断

  1. c.dataqsiz == 0代表为无缓冲的 channel,则判断当前 recvq 队列是否有等待接收的 goroutine,如果有则返回 false 代表可以非阻塞发送数据,否则返回 true 代表不可以非阻塞发送数据
  2. 如果是有缓冲的 channel 则判断 c.qcount == c.dataqsiz,如 c.qcount == c.dataqsiz 为 true 则代表缓冲区已满此时不可以非阻塞发送数据,否则返回 false 代表可以非阻塞发送数据

case2

send2
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
func chansend(c *hchan, ep unsafe.Pointer, block bool, callerpc uintptr) bool {
//..
lock(&c.lock)

if c.closed != 0 {
unlock(&c.lock)
panic(plainError("send on closed channel"))
}

if sg := c.recvq.dequeue(); sg != nil {
// Found a waiting receiver. We pass the value we want to send
// directly to the receiver, bypassing the channel buffer (if any).
send(c, sg, ep, func() { unlock(&c.lock) }, 3)
return true
}

//..
}

加锁完成以下操作

  1. 对已关闭的 channel 发送数据则会直接 panic
  2. 如果 recvq 队列中有等待读取的,则从中取出一个 goroutine 的封装对象 sudog,并通过 send 直接将元素拷贝交给 sudog 对应的 goroutine(基于 memmove 方法)
  3. 在 send 方法内会调用传入的函数解锁 channel 和 goready 唤醒对应的 goroutine,最终返回

case3

send3
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
func chansend(c *hchan, ep unsafe.Pointer, block bool, callerpc uintptr) bool {
//..
lock(&c.lock)
//..
if c.qcount < c.dataqsiz {
// Space is available in the channel buffer. Enqueue the element to send.
qp := chanbuf(c, c.sendx)
//..
typedmemmove(c.elemtype, qp, ep)
c.sendx++
if c.sendx == c.dataqsiz {
c.sendx = 0
}
c.qcount++
unlock(&c.lock)
return true
}

if !block {
unlock(&c.lock)
return false
}
//..
}

加锁完成以下操作

  1. 此时 channel 缓冲区有空间,无需阻塞写
  2. c.qcount < c.dataqsiz 为 true 则首先通过 chanbuf 方法获取到缓冲区写指针的地址 qp
  3. 调用 typedmemmove 方法将当前元素添加到环形缓冲区 qp 位置
  4. sendx++和 qcount++,解锁返回
  5. 如果前面操作均未完成,下面就要进入阻塞写流程,所以先判断是否要阻塞写,不阻塞写则直接解锁返回

case4

send4
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
func chansend(c *hchan, ep unsafe.Pointer, block bool, callerpc uintptr) bool {
//..
lock(&c.lock)
//..

gp := getg()
mysg := acquireSudog()
mysg.releasetime = 0
if t0 != 0 {
mysg.releasetime = -1
}

mysg.elem = ep
mysg.waitlink = nil
mysg.g = gp
mysg.isSelect = false
mysg.c = c
gp.waiting = mysg
gp.param = nil
c.sendq.enqueue(mysg)

gp.parkingOnChan.Store(true)
reason := waitReasonChanSend
if c.synctest {
reason = waitReasonSynctestChanSend
}
gopark(chanparkcommit, unsafe.Pointer(&c.lock), reason, traceBlockChanSend, 2)

KeepAlive(ep)

if mysg != gp.waiting {
throw("G waiting list is corrupted")
}
gp.waiting = nil
gp.activeStackChans = false
closed := !mysg.success
gp.param = nil
if mysg.releasetime > 0 {
blockevent(mysg.releasetime-t0, 2)
}
mysg.c = nil
releaseSudog(mysg)
if closed {
if c.closed == 0 {
throw("chansend: spurious wakeup")
}
panic(plainError("send on closed channel"))
}
return true
}

加锁完成以下操作

  1. 首先构造并封装 sudog,然后完成 sudog 和 goroutine 与 channel 的相关联
  2. 将 sudog 入队到当前 channel 的阻塞写协程队列,park 当前协程
  3. 如果当前协程被唤醒,首先判断 gp.waiting 是否被更改
  4. 然后 releaseSudog,然后判断 closed,如果是因为 chan 被关闭而醒来,就 panic;如果状态对不上,就认为出了 bug
  5. 解锁返回

读取数据

channel 读取数据有两种接收方式:

1
2
v := <-ch1
v, ok := <- ch1

使用第二种可以根据 bool 值判断,到底是读取到了零值还是 channel 已经 close

chanrecv1(c * hchan, elem unsafe.Pointer) And chanrecv2(c * hchan, elem unsafe.Pointer)

1
2
3
4
5
6
7
8
func chanrecv1(c *hchan, elem unsafe.Pointer) {
chanrecv(c, elem, true)
}

func chanrecv2(c *hchan, elem unsafe.Pointer) (received bool) {
_, received = chanrecv(c, elem, true)
return
}

前面两个读取数据的方法会被分别汇编为 chanrecv1 chanrecv2,不过核心还是 chanrecv 方法,下面我们展开阐述 chanrecv,在这里也分段阐述

case1

recv1 ****
1
2
3
4
5
6
7
8
9
10
11
func chanrecv(c *hchan, ep unsafe.Pointer, block bool) (selected, received bool) {
//..
if c == nil {
if !block {
return
}
gopark(nil, nil, waitReasonChanReceiveNilChan, traceBlockForever, 2)
throw("unreachable")
}
//..
}

block 参数标识是否以阻塞形式读取数据;如果向一个未初始化的 channel 读取数据并且 block 为 true,则 gopark 陷入死锁;否则直接返回 false

case2

recv2
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
func chanrecv(c *hchan, ep unsafe.Pointer, block bool) (selected, received bool) {
//..
if !block && empty(c) {
if atomic.Load(&c.closed) == 0 {
return
}
if empty(c) {
//..
if ep != nil {
typedmemclr(c.elemtype, ep)
}
return true, false
}
}
//..
}

快速处理路径:以非阻塞形式从 channel 中读取数据

  1. empty(c)函数会判断 channel 里是否有可读取的数据,此时有二种情况:1.无缓冲的 channel:查看是否有阻塞的等待写协程;2.有缓冲的 channel:查看缓冲区是否有数据;没有的话返回 true
  2. 进入 if,此时如果 channel 不是关闭的话直接返回
  3. 如果此时 channel 已关闭,再次判断看 empty 看缓冲区是否有残余数据,没有的话把接收缓冲 ep 清成该类型的零值,并返回 false

case3

recv3
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
func chanrecv(c *hchan, ep unsafe.Pointer, block bool) (selected, received bool) {
//..
lock(&c.lock)

if c.closed != 0 {
if c.qcount == 0 {
//..
unlock(&c.lock)
if ep != nil {
typedmemclr(c.elemtype, ep)
}
return true, false
}
} else {
if sg := c.sendq.dequeue(); sg != nil {
recv(c, sg, ep, func() { unlock(&c.lock) }, 3)
return true, true
}
}
//..
}

加锁完成以下操作

  1. 如果 channel 此时已关闭,先判断缓冲区是否还有数据,没有的话解锁把接收缓冲 ep 清成该类型的零值,并返回 false
  2. 此时 channel 还未关闭,并且有阻塞写协程,则直接调用 recv
  3. recv 处理有两种情况:1.channel 是无缓冲的,则直接读取写协程元素,并唤醒写协程,然后解锁返回;2.
    channel 是有缓冲的,则将缓冲区头部数据读取给读协程,然后将写协程元素放进缓冲区尾部,并唤醒写协程,然后解锁返回

case4

recv4
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
func chanrecv(c *hchan, ep unsafe.Pointer, block bool) (selected, received bool) {
//..
lock(&c.lock)
//..
if c.qcount > 0 {
// Receive directly from queue
qp := chanbuf(c, c.recvx)
//..
if ep != nil {
typedmemmove(c.elemtype, ep, qp)
}
typedmemclr(c.elemtype, qp)
c.recvx++
if c.recvx == c.dataqsiz {
c.recvx = 0
}
c.qcount--
unlock(&c.lock)
return true, true
}

if !block {
unlock(&c.lock)
return false, false
}
//..
}

加锁完成以下操作

  1. 此时 channel 缓冲区有数据,无阻塞写协程
  2. 直接从缓冲区 recvx index 对应位置的元素,赋值给 ep
  3. c.recvx++ 和 c.qcount–,解锁返回
  4. 前面流程都有没处理,此时缓冲区无数据且无阻塞写协程,所以进行是否阻塞判断,不阻塞等待的话直接解锁返回

case5

recv5
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
func chanrecv(c *hchan, ep unsafe.Pointer, block bool) (selected, received bool) {
//..
lock(&c.lock)
//..
gp := getg()
mysg := acquireSudog()
//..
mysg.elem = ep
mysg.waitlink = nil
gp.waiting = mysg

mysg.g = gp
mysg.isSelect = false
mysg.c = c
gp.param = nil
c.recvq.enqueue(mysg)
//..

gp.parkingOnChan.Store(true)
reason := waitReasonChanReceive
if c.synctest {
reason = waitReasonSynctestChanReceive
}
gopark(chanparkcommit, unsafe.Pointer(&c.lock), reason, traceBlockChanRecv, 2)

if mysg != gp.waiting {
throw("G waiting list is corrupted")
}
//..
gp.waiting = nil
gp.activeStackChans = false
if mysg.releasetime > 0 {
blockevent(mysg.releasetime-t0, 2)
}
success := mysg.success
gp.param = nil
mysg.c = nil
releaseSudog(mysg)
return true, success
}

加锁完成以下操作

  1. 获取到当前 goroutine 和构建一个封装的 sudog,然后完成 sudog 和 goroutine 与 channel 的相关联
  2. 将 sudog 入队到当前 channel 的阻塞读协程队列,park 当前协程
  3. 如果当前协程被唤醒,首先判断 gp.waiting 是否被更改
  4. 然后 releaseSudog,然后判断 closed,如果是因为 chan 被关闭而醒来,就返回零值并且 ok=false;如果状态对不上,就认为出了 bug
  5. 最后解锁返回

channel 的关闭

close

closechan(c * hchan)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
func closechan(c *hchan) {
if c == nil {
panic(plainError("close of nil channel"))
}

lock(&c.lock)
if c.closed != 0 {
unlock(&c.lock)
panic(plainError("close of closed channel"))
}
//。。
c.closed = 1
var glist gList
// release all readers
for {
sg := c.recvq.dequeue()
if sg == nil {
break
}
if sg.elem != nil {
typedmemclr(c.elemtype, sg.elem)
sg.elem = nil
}
if sg.releasetime != 0 {
sg.releasetime = cputicks()
}
gp := sg.g
gp.param = unsafe.Pointer(sg)
sg.success = false
if raceenabled {
raceacquireg(gp, c.raceaddr())
}
glist.push(gp)
}

// release all writers (they will panic)
for {
sg := c.sendq.dequeue()
if sg == nil {
break
}
sg.elem = nil
if sg.releasetime != 0 {
sg.releasetime = cputicks()
}
gp := sg.g
gp.param = unsafe.Pointer(sg)
sg.success = false
if raceenabled {
raceacquireg(gp, c.raceaddr())
}
glist.push(gp)
}
unlock(&c.lock)
// Ready all Gs now that we've dropped the channel lock.
for !glist.empty() {
gp := glist.pop()
gp.schedlink = 0
goready(gp, 3)
}
}
  1. 关闭一个未初始化的 channel 会直接 panic
  2. 加锁
  3. 关闭一个已经关闭的 channel,解锁,panic
  4. 将 closed 标识设为1
  5. 将当前 channel 的阻塞读/写队列中的协程全部加入到 glist
  6. 解锁,循环唤醒 glist 中阻塞的协程
  7. 返回

谢谢阅读!