channel
# 1.接口
# 1.1.实现方式
鸭子类型是动态编程语言的一种对象推断策略,它更关注对象能如何被使用,而不是对象的类型本身。
Go作为一种静态语言,它通过接口的方式完美支持鸭子类型。Go作为静态静态语言,不需要显示声明实现某个接口,编译器就会进行类型检查,只要实现了相关方法编译器就能检测,不需要像动态语言一样显示声明接口实现,运行期间才会检测到类型错误。type IGreeting interface { sayHello() } func sayHello(i IGreeting) { i.sayHello() } type Go struct {} func (g Go) sayHello() { fmt.Println("Hi, I am GO!") } type PHP struct {} func (p PHP) sayHello() { fmt.Println("Hi, I am PHP!") }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17鸭子类型是动态语言的一种风格,该风格指定一个对象有效的语义,不是由继承自特定的类或实现特定的接口,而是由它“当前方法和属性的集合”决定。Go作为一种静态语言,通过接口实现了鸭子类型,编译器隐匿了转换工作。
# 1.2.值接收和指针接收
方法可以对类型添加新的行为,与函数相比方法必须有一个接收者,接收者可以是
值类型,也可以是指针类型。调用方法时,不同类型接收者可以调用对方方法。package main import "fmt" type Person struct { age int } func (p Person) howOld() int { return p.age } func (p *Person) growUp() { p.age += 1 } func main() { // qcrao 是值类型 qcrao := Person{age: 18} // 值类型 调用接收者也是值类型的方法 fmt.Println(qcrao.howOld()) // 值类型 调用接收者是指针类型的方法 qcrao.growUp() fmt.Println(qcrao.howOld()) // ---------------------- // stefno 是指针类型 stefno := &Person{age: 100} // 指针类型 调用接收者是值类型的方法 fmt.Println(stefno.howOld()) // 指针类型 调用接收者也是指针类型的方法 stefno.growUp() fmt.Println(stefno.howOld()) }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实际上,当类型和方法接收者的类型不同时,编译器在背后作了一些工作,如表格描述所示:
值接收者 指针接收者 值类型调用者 方法会使用调用者的一个副本,类似于“传值” 使用值的引用来调用方法,上例中, qcrao.growUp()实际上是(&qcrao).growUp()指针类型调用者 指针被解引用为值,上例中, stefno.howOld()实际上是(*stefno).howOld()实际上也是“传值”,方法里的操作会影响到调用者,类似于指针传参,拷贝了一份指针 方法接收者是值类型还是指针类型,都可以通过值类型或指针类型调用,内部其实基于语法糖起作用。其实,实现了接收者是值类型的方法,相当于自动实现了接收者是指针类型的方法;实现了接收者是指针类型的方法,不会自动生成对应接收者是值类型的方法。
package main import "fmt" type coder interface { code() debug() } type Gopher struct { language string } func (p Gopher) code() { fmt.Printf("I am coding %s language\n", p.language) } func (p *Gopher) debug() { fmt.Printf("I am debuging %s language\n", p.language) } func main() { var c coder = &Gopher{"Go"} c.code() c.debug() }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上述程序运行正常,但如果把初始化时的指针去掉,运行会报错:
src/main.go:23:6: cannot use Gopher literal (type Gopher) as type coder in assignment: Gopher does not implement coder (debug method has pointer receiver)1
2报错会提示
Gopher没有实现coder,为什么?前边我们提过,接收者是值类型相当于具备了指针类型方法,所以*Gopher其实实现了coder,但接收者是指针类型不会自动实现值类型的方法,所以Gopher只实现了coder的code方法,无法赋值给coder作为实现。Go语言中,接收者使用值类型还是指针类型由内部成员决定,如果成员是内置的基本类型,则可以使用值类型接收者;如果内部成员是引用类型,,包括slice、map、interface、channel,或者无法被安全复制的共享成员struct file,接收者使用指针类型可以复制指针,保证成员内存区域只有一份,避免方法调用时拷贝成员内容。
# 1.3.iface和eface
iface和eface是Go中描述接口的底层结构体,区别是iface描述的接口包含方法,eface是不包含任何方法的空接口interface{}。type iface struct { tab *itab data unsafe.Pointer } type itab struct { inter *interfacetype // 描述接口类型 _type *_type // 实体类型,包括内存对齐、大小等 hash uint32 // copy of _type.hash. Used for type switches. _ [4]byte fun [1]uintptr // 放置和接口方法对应的具体数据类型的方法地址,实现接口方法调用的动态分派。 }1
2
3
4
5
6
7
8
9
10
11
12iface内部维护两个指针,tab指向一个itab实体,它表示接口的类型以及赋给这个接口的实体类型,data指向接口具体值,一般是一个指向堆内存的指针。itab一般在每次给接口赋值发生转换时更新,或者直接从缓存获取itab。itab中fun字段存储方法的函数指针时,会将方法按照名称的字典序存到连续内存空间。从汇编角度来看,通过增加地址就能获取到这些函数指针。type interfacetype struct { typ _type pkgpath name mhdr []imethod }1
2
3
4
5interfacetype描述接口类型,内部包装_type,_type实际上描述Go语言中各种数据类型的结构体,mhdr表示接口所定义的函数列表,pkgpath记录定义了接口的包名。
type eface struct { _type *_type data unsafe.Pointer }1
2
3
4相比
iface,eface相对简单,只维护了一个_type字段,表示空接口所承载的具体的实体类型,data描述了具体的值。
type _type struct { // 类型大小 size uintptr ptrdata uintptr // 类型的 hash 值 hash uint32 // 类型的 flag,和反射相关 tflag tflag // 内存对齐相关 align uint8 fieldalign uint8 // 类型的编号,有bool, slice, struct 等等等等 kind uint8 alg *typeAlg // gc 相关 gcdata *byte str nameOff ptrToThis typeOff }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19Go语言各种数据类型都是在_type字段的基础上,增加一些额外的字段进行管理。type arraytype struct { typ _type elem *_type slice *_type len uintptr } type chantype struct { typ _type elem *_type dir uintptr } type slicetype struct { typ _type elem *_type } type structtype struct { typ _type pkgPath name fields []structfield }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
# 1.4.接口的动态类型和动态值
iface包含两个字段,tab是接口表指针,指向类型信息;data是数据指针,指向值信息,它们分别被称为动态类型和动态值,接口值包括动态类型和动态值。package main import "fmt" type Coder interface { code() } type Gopher struct { name string } func (g Gopher) code() { fmt.Printf("%s is coding\n", g.name) } func main() { var c Coder fmt.Println(c == nil) fmt.Printf("c: %T, %v\n", c, c) var g *Gopher fmt.Println(g == nil) c = g fmt.Println(c == nil) fmt.Printf("c: %T, %v\n", c, c) } --- 输出 true c: <nil>, <nil> true false c: *main.Gopher, <nil>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接口类型的零值指的是
动态类型和动态值都是nil,当两部分的值都是nil的情况下,这个接口值才会被认为是nil。c的动态类型和动态值都是nil,g也是nil,当把g赋值给c后,c的动态类型变成了*main.Gopher,仅管c的动态值仍为nil,但是当c和nil比较时候,结果就是false。package main import ( "unsafe" "fmt" ) type iface struct { itab, data uintptr } func main() { var a interface{} = nil var b interface{} = (*int)(nil) x := 5 var c interface{} = (*int)(&x) ia := *(*iface)(unsafe.Pointer(&a)) ib := *(*iface)(unsafe.Pointer(&b)) ic := *(*iface)(unsafe.Pointer(&c)) fmt.Println(ia, ib, ic) fmt.Println(*(*int)(unsafe.Pointer(ic.data))) }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程序中定义了一个
iface结构体,用两个指针来描述itab和data,之后将a b c在内存中的内容强制解释成自定义的iface,最后就可以打印出动态类型和动态值的地址。{0 0} {17426912 0} {17426912 842350714568} 51
2a的动态类型和动态值的地址均为0,也就是nil;b的动态类型和c的动态类型一致,都是*int。
# 1.5.接口实现检测
很多情况下,一些开源库会用强制转换判断类型是否实现接口,如下:
var _ io.Writer = (*myWriter)(nil)1实际上这种方式会经过编译器检测,判断
*myWriter类型是否实现了io.Writer接口。package main import "io" type myWriter struct { } /*func (w myWriter) Write(p []byte) (n int, err error) { return }*/ func main() { // 检查 *myWriter 类型是否实现了 io.Writer 接口 var _ io.Writer = (*myWriter)(nil) // 检查 myWriter 类型是否实现了 io.Writer 接口 var _ io.Writer = myWriter{} }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19实际上,上述赋值语句会发生隐式地类型转换,转换过程中,编译器会检测等号右边的类型是否实现了等号左边接口规定的函数。
# 1.6.接口构造
package main import "fmt" type Person interface { growUp() } type Student struct { age int } func (p Student) growUp() { p.age += 1 return } func main() { var qcrao = Person(Student{age: 18}) fmt.Println(qcrao) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22执行
go tool compile -S main.go编译命令后可以得到main函数的汇编代码。--- main.Student.growUp函数的汇编实现 0x0000 00000 (main.go:24) TEXT main.Student.growUp(SB), NOSPLIT|ABIInternal, $0-8 0x0000 00000 (main.go:24) FUNCDATA $0, gclocals·g2BeySu+wFnoycgXfElmcg==(SB) 0x0000 00000 (main.go:24) FUNCDATA $1, gclocals·g2BeySu+wFnoycgXfElmcg==(SB) 0x0000 00000 (main.go:24) FUNCDATA $5, main.Student.growUp.arginfo1(SB) 0x0000 00000 (main.go:24) FUNCDATA $6, main.Student.growUp.argliveinfo(SB) 0x0000 00000 (main.go:24) PCDATA $3, $1 0x0000 00000 (main.go:26) RET --- main函数内容 . main.main STEXT size=111 args=0x0 locals=0x48 funcid=0x0 align=0x0 0x0000 00000 (main.go:29) TEXT main.main(SB), ABIInternal, $72-0 0x0000 00000 (main.go:29) CMPQ SP, 16(R14) 0x0004 00004 (main.go:29) PCDATA $0, $-2 0x0004 00004 (main.go:29) JLS 104 0x0006 00006 (main.go:29) PCDATA $0, $-1 0x0006 00006 (main.go:29) SUBQ $72, SP 0x000a 00010 (main.go:29) MOVQ BP, 64(SP) 0x000f 00015 (main.go:29) LEAQ 64(SP), BP --- 接口构造,用于将main.Student类型的数据和类型信息封装成一个接口值,这是Go处理接口的标准形式 0x0014 00020 (main.go:29) FUNCDATA $0, gclocals·g2BeySu+wFnoycgXfElmcg==(SB) 0x0014 00020 (main.go:29) FUNCDATA $1, gclocals·EaPwxsZ75yY1hHMVZLmk6g==(SB) 0x0014 00020 (main.go:29) FUNCDATA $2, main.main.stkobj(SB) // 将立即数18存入栈上的某个位置,用于后续可能的接口类型构造 0x0014 00020 (main.go:30) MOVQ $18, main..autotmp_9+40(SP) // 将栈上的main..autotmp_9+40(SP)处的值加载到AX寄存器中,作为接口的具体数据值 0x001d 00029 (main.go:30) MOVQ main..autotmp_9+40(SP), AX 0x0022 00034 (main.go:30) PCDATA $1, $0 // 用于转换类型和类型检查 0x0022 00034 (main.go:30) CALL runtime.convT64(SB) // 将X15寄存器的值存储到栈上某个位置 0x0027 00039 (main.go:32) MOVUPS X15, main..autotmp_13+48(SP) // 将main.Student类型的接口表格信息存储到CX寄存器中 0x002d 00045 (main.go:32) MOVQ go.itab.main.Student,main.Person+8(SB), CX // 将CX寄存器的值存储到栈上,包括接口类型信息 0x0034 00052 (main.go:32) MOVQ CX, main..autotmp_13+48(SP) // 将AX寄存器中的具体数据存储到栈上 0x0039 00057 (main.go:32) MOVQ AX, main..autotmp_13+56(SP) --- 输出接口值 0x003e 00062 (<unknown line number>) NOP // 加载标准输出os.Stdout的地址到BX寄存器 0x003e 00062 ($GOROOT/src/fmt/print.go:294) MOVQ os.Stdout(SB), BX // 将os.File类型的接口信息加载到AX寄存器 0x0045 00069 ($GOROOT/src/fmt/print.go:294) LEAQ go.itab.*os.File,io.Writer(SB), AX // 将之前构造的接口值的地址加载到CX寄存器 0x004c 00076 ($GOROOT/src/fmt/print.go:294) LEAQ main..autotmp_13+48(SP), CX // 设置参数1,用于表示fmt.Fprintln的第一个参数是os.Stdout 0x0051 00081 ($GOROOT/src/fmt/print.go:294) MOVL $1, DI // 将DI的值移动到SI寄存器 0x0056 00086 ($GOROOT/src/fmt/print.go:294) MOVQ DI, SI 0x0059 00089 ($GOROOT/src/fmt/print.go:294) CALL fmt.Fprintln(SB)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
5310~14行会调用runtime.convT64(SB)函数,这个函数的形式如下:func convT64(val uint64) (x unsafe.Pointer) { ... TODO }1
2
3convT64会构造出一个unsafe.Pointer,参数位置是SB,这里被赋值go.itab."".Student,"".Person(SB)的地址。go.itab.main.Student,main.Person SRODATA dupok size=32 0x0000 00 00 00 00 00 00 00 00 00 00 00 00 00 00 00 00 ................ 0x0010 eb 29 b9 7a 00 00 00 00 00 00 00 00 00 00 00 00 .).z............ rel 0+8 t=1 type.main.Person+0 rel 8+8 t=1 type.main.Student+0 rel 24+8 t=-32767 main.(*Student).growUp+01
2
3
4
5
6这里
size大小是32字节,可以回顾一下itab表字段。type itab struct { inter *interfacetype // 8字节 _type *_type // 8字节 hash uint32 // 4字节 _ [4]byte // 4字节 fun [1]uintptr // 8字节 }1
2
3
4
5
6
7每个字段的大小相加,
itab结构体的大小就是32字节,上面那一串数字其实是itab序列化后的内容,注意到大部分数字为0,从17开始的4字节其实是itab的哈希值,用于判断两个类型是否相同。下面两行是链接指令,用于将所有源文件综合起来,给每个符号赋予一个全局的位置值。这里其实比较明确,前8字节存储type."".Person的地址,对应itab里的inter字段,表示接口类型;8~16字节存储type."".Student的地址,对应itab里的_type字段,表示具体类型;最后的链接对应24开始的8字节,存储类型实现的方法地址,对应fun字段。第二个参数就比较简单,它代表值的地址,这里是18的地址,也就是初始化Student结构体时指定的字段值。func convT64(val uint64) (x unsafe.Pointer) { if val < uint64(len(staticuint64s)) { x = unsafe.Pointer(&staticuint64s[val]) } else { x = mallocgc(8, uint64Type, false) *(*uint64)(x) = val } return }1
2
3
4
5
6
7
8
9
# 1.7.类型转换和断言
Go语言中不允许隐式类型转换,也就是=两边不允许出现类型不相同的变量。类型转换、类型断言本质都是把一个类型转换成另一个类型。不同在于,类型断言针对接口变量进行操作。func main() { var i int = 9 var f float64 f = float64(i) fmt.Printf("%T, %v\n", f, f) f = 10.8 a := int(f) fmt.Printf("%T, %v\n", a, a) // s := []int(i) }1
2
3
4
5
6
7
8
9
10
11
12
13对于
类型转换而言,转换前后的两个类型需要相互兼容,如int和float64类型。但对于类型断言来说,主要用于接口真实类型判断,即空接口类型断言为具体的类型。type Student struct { Name string Age int } func main() { var i interface{} = new(Student) s := i.(*Student) fmt.Println(s) }1
2
3
4
5
6
7
8
9
10
11特别地,
fmt.Println函数的参数是interface,对于内置类型会使用穷举获取它的真实类型,然后转换为字符串打印;对于自定义类型,会确认该类型是否实现String()方法,实现会直接打印String()方法结果,否则利用反射遍历对象的成员进行打印。
# 1.8.接口转换原理
iface实际上包含接口的类型interfacetype和实体类型的类型_type,这两者都是iface的字段itab的成员,也就是说生成一个itab同时需要接口类型和实体类型。当判断一种类型是否满足某个接口时,Go使用类型的方法集和接口所需要的方法集进行匹配,如果类型的方法集完全包含接口的方法集,可以认为类型实现了接口。方法集进行匹配时,会对方法集进行名字的字典序排序,所以实际上只需要O(m+n)时间复杂度。package main import "fmt" type coder interface { code() run() } type runner interface { run() } type Gopher struct { language string } func (g Gopher) code() { return } func (g Gopher) run() { return } func main() { var c coder = Gopher{} var r runner r = c fmt.Println(c, r) }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上述程序定义了两个
interface,分别是coder和runner,定义了一个实体类型Gopher,类型Gopher实现了两个方法,分别是run和code。main函数里定义了一个接口变量c,绑定了一个Gopher对象,之后赋值给另一个接口变量r,这是因为c中包含run方法,这样就实现了接口转换。执行编译命令后,可以看到接口转换对应的汇编代码,其中包括一个关键函数:0x0025 00037 (main.go:36) CALL runtime.convTstring(SB) 0x002a 00042 (main.go:36) MOVQ AX, main..autotmp_27+40(SP) 0x002f 00047 (main.go:39) LEAQ go.itab.main.Gopher,main.coder(SB), BX 0x0036 00054 (main.go:39) LEAQ type.main.runner(SB), AX 0x003d 00061 (main.go:39) PCDATA $1, $1 0x003d 00061 (main.go:39) NOP 0x0040 00064 (main.go:39) CALL runtime.convI2I(SB) 0x0045 00069 (main.go:40) MOVUPS X15, main..autotmp_15+64(SP)1
2
3
4
5
6
7
8runtime.convI2I(SB)调用其实就对应了r=c,这个函数主要作用就是把一个接口的itab转换成另一个接口的itab。func convI2I(dst *interfacetype, src *itab) *itab { if src == nil { return nil } if src.inter == dst { return src } return getitab(dst, src._type, false) }1
2
3
4
5
6
7
8
9程序实现比较简单,
dst入参表示接口类型,src表示绑定了实体类型的接口itab。通过前面的分析,iface是由tab和data两个字段组成。所以,实际上该函数就是要把itab转换成新接口的itab。其中,tab是由接口类型interfacetype和实体类型_type组成,所以最关键的就是getitab(dst, src._type, false)。func getitab(inter *interfacetype, typ *_type, canfail bool) *itab { // 空接口没有方法表,抛出错误 if len(inter.mhdr) == 0 { throw("internal error - misuse of itab") } ... var m *itab // 使用原子操作从itabtable中查找itab t := (*itabTableType)(atomic.Loadp(unsafe.Pointer(&itabTable))) // 找到则返回 if m = t.find(inter, typ); m != nil { goto finish } // 没找到时,获取锁再查找一次 lock(&itabLock) if m = itabTable.find(inter, typ); m != nil { unlock(&itabLock) goto finish } // itab不存在时,分配内存并初始化新的itab m = (*itab)(persistentalloc(unsafe.Sizeof(itab{})+uintptr(len(inter.mhdr)-1)*goarch.PtrSize, 0, &memstats.other_sys)) // 指定接口类型 m.inter = inter // 指定实体类型 m._type = typ // 初始化其他参数 m.hash = 0 m.init() // 添加到itab表 itabAdd(m) unlock(&itabLock) finish: // 实体实现了接口方法正常返回 if m.fun[0] != 0 { return m } // 抛出异常或结束 if canfail { return nil } panic(&TypeAssertionError{concrete: typ, asserted: &inter.typ, missingMethod: m.init()}) } func (t *itabTableType) find(inter *interfacetype, typ *_type) *itab { mask := t.size - 1 // 计算哈希值 h := itabHashFunc(inter, typ) & mask // 使用二次探测法解决哈希冲突查找itab for i := uintptr(1); ; i++ { p := (**itab)(add(unsafe.Pointer(&t.entries), h*goarch.PtrSize)) // 基于哈希偏移原子查找itab m := (*itab)(atomic.Loadp(unsafe.Pointer(p))) // 没找到提前结束 if m == nil { return nil } // 找到返回对应itab地址 if m.inter == inter && m._type == typ { return m } h += i h &= mask } } func itabAdd(m *itab) { // 判断当前是否执行内存分配,正在分配则抛出异常拒绝访问 if getg().m.mallocing != 0 { throw("malloc deadlock") } t := itabTable // 哈希表扩容条件判断,超出0.75负载因子则进行扩容,减少槽位冲突 if t.count >= 3*(t.size/4) { // 75% load factor // 扩容哈希表为2倍 t2 := (*itabTableType)(mallocgc((2+2*t.size)*goarch.PtrSize, nil, true)) t2.size = t.size * 2 // 复制旧的哈希表 iterate_itabs(t2.add) if t2.count != t.count { throw("mismatched count during itab table copy") } // 原子操作更新哈希表 atomicstorep(unsafe.Pointer(&itabTable), unsafe.Pointer(t2)) t = itabTable } // 插入新的itab条目 t.add(m) }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
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95总结来说,
getitab函数会根据interfacetype和_type去全局的itab哈希表中查找,找到直接返回,否则根据给定的接口类型和实体类型新生成一个itab插入到itab哈希表,这样下一次就可以直接拿到itab。当把实体类型赋值给接口的时候,会调用conv系列函数,例如空接口调用convT2E系列,非空接口调用convT2I系列,这些函数比较类型:- 具体类型转空接口时,
_type字段直接复制源类型的_type;调用mallocgc获得一块新内存,把值赋值进去,data指向这块新内存 - 具体类型转非空接口时,入参
tab是编译器在编译阶段预先生成好的,新接口tab字段直接指向入参tab指向的itab;调用mallocgc获得一块新内存,把值复制进去,data再指向这块新内存 - 对于接口转接口,
itab调用getitab函数获取,只用生成一次,之后直接从hash表中获取
- 具体类型转空接口时,
# 2.通道
# 2.1.底层结构
type hchan struct { // chan里元素数量 qcount uint // chan底层循环数组的长度 dataqsiz uint // 指向底层循环数组的指针 // 只针对有缓冲的channel buf unsafe.Pointer // chan中元素大小 elemsize uint16 // chan是否被关闭的标志 closed uint32 // chan中元素类型 elemtype *_type // element type // 已发送元素在循环数组中的索引 sendx uint // send index // 已接收元素在循环数组中的索引 recvx uint // receive index // 等待接收的 goroutine 队列 recvq waitq // list of recv waiters // 等待发送的 goroutine 队列 sendq waitq // list of send waiters // 保护 hchan 中所有字段 lock mutex }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25channel定义中,buf指向底层循环数组,只有缓冲型的chan才有,sendx和recvx均指向底层循环数组,表示当前可以发送和接收的元素位置索引值,sendq和recvq分别表示被阻塞的goroutine,这些goroutine由于尝试读取channel或向channel发送数据而被阻塞。waitq是sudog的一个双向链表,sudog是对goroutine的一个封装:type waitq struct { first *sudog last *sudog }1
2
3
4lock用来保证每个读channel或写channel的操作都是原子的。例如容量为6,元素为int型的channel数据结构如下:
通道的创建一般会使用
make实现,底层其实会调用makechan:func makechan(t *chantype, size int) *hchan { ... }1
2
3从函数原型看,创建的
chan是一个指针,所以chan可以在函数直接传递,而不用显示声明为指针。const hchanSize = unsafe.Sizeof(hchan{}) + uintptr(-int(unsafe.Sizeof(hchan{}))&(maxAlign-1)) func makechan(t *chantype, size int64) *hchan { elem := t.elem ... var c *hchan // 如果元素类型不含指针或者size大小为0(无缓冲类型) // 只进行一次内存分配 if elem.kind&kindNoPointers != 0 || size == 0 { // 如果hchan结构体中不含指针,GC就不会扫描chan中的元素 // 只分配 "hchan结构体大小+元素大小*个数"的内存 c = (*hchan)(mallocgc(hchanSize+uintptr(size)*elem.size, nil, true)) // 如果是缓冲型channel且元素大小不等于0 if size > 0 && elem.size != 0 { c.buf = add(unsafe.Pointer(c), hchanSize) } else { // 1. 非缓冲型的,buf没用,直接指向chan起始地址处 // 2. 缓冲型的,元素无指针且元素类型为struct{},只会用到接收和发送游标,不会真正拷贝东西到c.buf处 c.buf = unsafe.Pointer(c) } } else { // 进行两次内存分配操作 c = new(hchan) c.buf = newarray(elem, int(size)) } c.elemsize = uint16(elem.size) c.elemtype = elem // 循环数组长度 c.dataqsiz = uint(size) // 返回hchan指针 return c }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
# 2.2.chan发送
发送操作最终转化为
chansend函数,位于src/runtime/chan.go。func chansend(c *hchan, ep unsafe.Pointer, block bool, callerpc uintptr) bool { // 如果channel是nil if c == nil { // 不能阻塞,直接返回false,表示未发送成功 if !block { return false } // 当前goroutine被挂起 gopark(nil, nil, "chan send (nil chan)", traceEvGoStop, 2) throw("unreachable") } ... // 对于不阻塞的 send,快速检测失败场景 // 如果channel未关闭且channel没有多余的缓冲空间 // 1. channel是非缓冲型的,且等待接收队列里没有goroutine // 2. channel是缓冲型的,但循环数组已经装满了元素 if !block && c.closed == 0 && ((c.dataqsiz == 0 && c.recvq.first == nil) || (c.dataqsiz > 0 && c.qcount == c.dataqsiz)) { return false } var t0 int64 if blockprofilerate > 0 { t0 = cputicks() } // 锁住channel,并发安全 lock(&c.lock) // 如果channel关闭了 if c.closed != 0 { // 解锁 unlock(&c.lock) // 直接 panic panic(plainError("send on closed channel")) } // 如果接收队列里有goroutine,直接将要发送的数据拷贝到接收goroutine if sg := c.recvq.dequeue(); sg != nil { send(c, sg, ep, func() { unlock(&c.lock) }, 3) return true } // 对于缓冲型的channel,如果还有缓冲空间 if c.qcount < c.dataqsiz { // qp指向buf的sendx 位置 qp := chanbuf(c, c.sendx) ... // 将数据从ep处拷贝到qp typedmemmove(c.elemtype, qp, ep) // 发送游标值加1 c.sendx++ // 如果发送游标值等于容量值,游标值归0 if c.sendx == c.dataqsiz { c.sendx = 0 } // 缓冲区的元素数量加一 c.qcount++ // 解锁 unlock(&c.lock) return true } // 如果不需要阻塞,则直接返回错误 if !block { unlock(&c.lock) return false } // channel满,发送方会被阻塞,接下来会构造一个sudog // 获取当前goroutine的指针 gp := getg() mysg := acquireSudog() mysg.releasetime = 0 if t0 != 0 { mysg.releasetime = -1 } mysg.elem = ep mysg.waitlink = nil mysg.g = gp mysg.selectdone = nil mysg.c = c gp.waiting = mysg gp.param = nil // 当前goroutine进入发送等待队列 c.sendq.enqueue(mysg) // 当前goroutine被挂起 goparkunlock(&c.lock, "chan send", traceEvGoBlockSend, 3) // 从这里开始被唤醒(channel有机会可以发送) if mysg != gp.waiting { throw("G waiting list is corrupted") } gp.waiting = nil if gp.param == nil { if c.closed == 0 { throw("chansend: spurious wakeup") } // 被唤醒后,channel关闭抛panic panic(plainError("send on closed channel")) } gp.param = nil if mysg.releasetime > 0 { blockevent(mysg.releasetime-t0, 2) } // 去掉mys上绑定的channel mysg.c = nil releaseSudog(mysg) return true }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
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118在不阻塞的场景下快速检测发送失败主要是以下程序:
if !block && c.closed == 0 && ((c.dataqsiz == 0 && c.recvq.first == nil) || (c.dataqsiz > 0 && c.qcount == c.dataqsiz)) { return false }1
2
3if条件会先读取两个变量block和c.closed,block作为函数的参数不会变,c.closed可能被其他协程改变。c.dataqsize创建时已经确定,也不会被修改,这里不加锁影响的是c.qcount和c.recvq.first。这里其实就是读两个word操作,非缓冲型场景的c.closed和c.recvq.first和缓冲型场景的c.qcount。当c.closed=0为真,观测到c.recvq.first==nil或者c.qcount=d.dataqsize,说明这次发送操作失败。该断言涉及两个观测项,channel未关闭、channel not ready for sending,这两项因此没加锁可能出现观测前后不一致情况。但是,因为一个closed channel不能将channel状态从ready for sending变成not ready for sending。所以当检测到not ready for sending时,即使channel在某个观测中间被关闭,那也说明在两个观测中间channel满足not closed和not ready for sending,快速失败返回false也不会有问题。如果可以从等待接收队列
recvq出队一个sudog,说明此时chan是空的,没有元素。这是会调用send函数直接从发送者的栈拷贝到接收者的栈,关键操作由sendDirect函数完成。// send函数处理向一个空的channel发送操作 // ep指向被发送的元素,会被直接拷贝到接收的goroutine,然后接收的goroutine会被唤醒 // c必须是空的(因为等待队列里有goroutine,肯定是空的) // c必须被上锁,发送操作执行完后,会使用unlockf函数解锁 // sg必须已经从等待队列里取出来 // ep必须是非空,并且它指向堆或调用者的栈 func send(c *hchan, sg *sudog, ep unsafe.Pointer, unlockf func(), skip int) { ... // sg.elem指向接收到的值存放的位置,如 val <- ch,指的就是 &val if sg.elem != nil { // 直接拷贝内存(从发送者到接收者) sendDirect(c.elemtype, sg, ep) sg.elem = nil } // sudog上绑定的goroutine gp := sg.g // 解锁 unlockf() gp.param = unsafe.Pointer(sg) if sg.releasetime != 0 { sg.releasetime = cputicks() } // 唤醒接收的goroutine.skip和打印栈相关,暂时不理会 goready(gp, skip+1) } // 向一个非缓冲型的channel发送数据、从一个无元素的(非缓冲型或缓冲型但空)的channel接收数据 // 都会导致一个goroutine直接操作另一个goroutine的栈 // 由于GC 假设对栈的写操作只能发生在goroutine正在运行中并且由当前goroutine来写 // 所以这里实际上违反了这个假设,可能会造成一些问题,所以需要用到写屏障来规避 func sendDirect(t *_type, sg *sudog, src unsafe.Pointer) { // src在当前goroutine的栈上,dst是另一个goroutine的栈 // 直接进行内存"搬迁" // 如果目标地址的栈发生了栈收缩,当我们读出了sg.elem后 // 就不能修改真正的dst位置的值了 // 因此需要在读和写之前加上一个屏障 dst := sg.elem typeBitsBulkBarrier(t, uintptr(dst), uintptr(src), t.size) memmove(dst, src, t.size) }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这里涉及到一个
goroutine直接写另一个goroutine栈的操作,不同协程的一般各自独有,这违反了GC的一些假设。基于此,写的过程中会增加写屏障,保证正确地完成写操作。这样的好处是减少了一次内存copy,不用先拷贝到channel.buf,直接由发送者到接收者,提高效率。然后,解锁、唤醒接收者,等待调度器光临,接收者继续执行接收操作之后的逻辑。如果
c.qcount<c.dataqsize,说明缓冲区可用,先通过函数取出待发送元素应该去到地位置。qp := chanbuf(c, c.sendx) // 返回循环队列里第 i 个元素的地址处 func chanbuf(c *hchan, i uint) unsafe.Pointer { return add(c.buf, uintptr(i)*uintptr(c.elemsize)) }1
2
3
4
5
6c.sendx指向下一个待发送元素在循环数组中的位置,然后调用typedmemmove函数将其拷贝到循环数组,之后c.sendx+1,元素总量加1,最后解锁并返回。如果没有命中以上条件,说明channel已经满了,不管这个channel是缓冲型还是非缓冲型,都要将这个sender关起来;如果真的阻塞,会构建sudog入队到c.sendq,调用goparkunlock将当前协程挂起、解锁,等待合适的机会唤醒。唤醒后,会有一些绑定操作,sudog通过g字段绑定goroutine,goroutine通过waiting绑定sudog,sudog通过elem字段绑定待发送元素地址,以及c字段绑定被坑在此处的channel。所以,待发送的元素地址其实存储在sudog结构体里,也就是当前goroutine里。func goroutineA(a <-chan int) { val := <- a fmt.Println("goroutine A received data: ", val) return } func goroutineB(b <-chan int) { val := <- b fmt.Println("goroutine B received data: ", val) return } func main() { ch := make(chan int) go goroutineA(ch) go goroutineB(ch) ch <- 3 time.Sleep(time.Second) ch1 := make(chan struct{}) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21上述程序中,
G1和G2会被挂起,等待sender的解救。当主协程向ch发送一个元素3时,会发现recvq队列有接收者等待,就会出队一个sudog,把recvq里first指针的元素推举出来,将其加入到P的可运行goroutine队列。sender会把发送元素拷贝到sudog的elem地址处,最后调用goready将G1唤醒,状态变为runnable。
当调度器光顾
G1时,将G1变为running状态,执行goroutineA剩余逻辑,G表示其他可能的goroutine。这里其实涉及到一个协程写另一个协程栈的操作,有两个receiver在channel的recvq中等待接收数据,sender准备向channel发送数据时,为了高效不会通过channel.buf多中转拷贝一次,直接从源地址把数据copy到目的地址。
# 2.3.chan接收
chan接收操作有两种写法,一种带ok,反应channel是否关闭,另一种不带ok单纯取值。其实调用函数经过编译器的处理后,对应源码里的两个函数:// entry points for <- c from compiled code 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 }1
2
3
4
5
6
7
8
9两个函数接受值都比较特殊,会放到参数
elem所指向的地址,如果忽略则elem=nil。两者最终都转向chanrecv函数。// chanrecv函数接收channel c的元素并将其写入ep所指向的内存地址 // 如果ep是nil,说明忽略接收值 // 如果block == false,即非阻塞型接收,在没有数据可接收的情况下返回(false,false) // 否则,如果c处于关闭状态,将ep指向的地址清零,返回(true,false) // 否则,用返回值填充ep指向的内存地址,返回(true, true) // 如果ep非空,则应该指向堆或者函数调用者的栈 func chanrecv(c *hchan, ep unsafe.Pointer, block bool) (selected, received bool) { ... // 如果是一个nil的channel if c == nil { // 如果不阻塞,直接返回(false, false) if !block { return } // 否则,接收一个nil的channel,goroutine挂起 gopark(nil, nil, "chan receive (nil chan)", traceEvGoStop, 2) // 不会执行到这里 throw("unreachable") } // 在非阻塞模式下,快速检测到失败,不用获取锁,快速返回 // 当我们观察到channel没准备好接收: // 1. 非缓冲型,等待发送列队sendq里没有goroutine在等待 // 2. 缓冲型,但buf里没有元素 // 之后,又观察到closed == 0,即channel未关闭。 // 因为channel不可能被重复打开,所以前一个观测的时候channel也是未关闭的 // 因此在这种情况下可以直接宣布接收失败,返回(false,false) if !block && (c.dataqsiz == 0 && c.sendq.first == nil || c.dataqsiz > 0 && atomic.Loaduint(&c.qcount) == 0) && atomic.Load(&c.closed) == 0 { return } var t0 int64 if blockprofilerate > 0 { t0 = cputicks() } // 加锁 lock(&c.lock) // channel已关闭,并且循环数组buf里没有元素 // 这里可以处理非缓冲型关闭和缓冲型关闭但buf无元素的情况 // 也就是说即使是关闭状态,但在缓冲型的channel,buf里有元素的情况下还能接收到元素 if c.closed != 0 && c.qcount == 0 { if raceenabled { raceacquire(unsafe.Pointer(c)) } // 解锁 unlock(&c.lock) if ep != nil { // 从一个已关闭的channel执行接收操作,且未忽略返回值 // 那么接收的值将是一个该类型的零值 // typedmemclr根据类型清理相应地址的内存 typedmemclr(c.elemtype, ep) } // 从一个已关闭的channel接收,selected会返回true return true, false } // 等待发送队列里有goroutine存在,说明buf是满的 // 这有可能是: // 1. 非缓冲型的channel // 2. 缓冲型的channel,但buf满了 // 针对 1,直接进行内存拷贝(从sender goroutine -> receiver goroutine) // 针对 2,接收到循环数组头部的元素,并将发送者的元素放到循环数组尾部 if sg := c.sendq.dequeue(); sg != nil { recv(c, sg, ep, func() { unlock(&c.lock) }, 3) return true, true } // 缓冲型,buf里有元素,可以正常接收 if c.qcount > 0 { // 直接从循环数组里找到要接收的元素 qp := chanbuf(c, c.recvx) ... // 代码里没有忽略要接收的值,不是"<- ch",而是"val <- ch",ep指向val if ep != nil { typedmemmove(c.elemtype, ep, qp) } // 清理掉循环数组里相应位置的值 typedmemclr(c.elemtype, qp) // 接收游标向前移动 c.recvx++ // 接收游标归零 if c.recvx == c.dataqsiz { c.recvx = 0 } // buf数组里的元素个数减1 c.qcount-- // 解锁 unlock(&c.lock) return true, true } if !block { // 非阻塞接收,解锁,selected返回false,因为没有接收到值 unlock(&c.lock) return false, false } // 接下来就是要被阻塞的情况了 // 构造一个sudog gp := getg() mysg := acquireSudog() mysg.releasetime = 0 if t0 != 0 { mysg.releasetime = -1 } // 待接收数据的地址保存下来 mysg.elem = ep mysg.waitlink = nil gp.waiting = mysg mysg.g = gp mysg.selectdone = nil mysg.c = c gp.param = nil // 进入channel的等待接收队列 c.recvq.enqueue(mysg) // 将当前goroutine挂起 goparkunlock(&c.lock, "chan receive", traceEvGoBlockRecv, 3) // 被唤醒,接着从这里继续执行一些扫尾工作 if mysg != gp.waiting { throw("G waiting list is corrupted") } gp.waiting = nil if mysg.releasetime > 0 { blockevent(mysg.releasetime-t0, 2) } closed := gp.param == nil gp.param = nil mysg.c = nil releaseSudog(mysg) return true, !closed }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
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141核心流程
1.
channel是一个空值,非阻塞模式下直接返回,阻塞模式下调用gopark函数挂起goroutine,这种情况下会一直阻塞下去2.非阻塞模式下快速检测失败,处理部分边界条件
3.加锁判断,如果
channel关闭,循环数组buf里没有元素,对应非缓冲型关闭和缓冲型关闭buf无元素情况,返回对应类型的零值,recived标识false,如果处于select语境,属于被命中情况4.如果有等待发送的队列,
channel已满,调用recv接收数据func recv(c *hchan, sg *sudog, ep unsafe.Pointer, unlockf func(), skip int) { // 如果是非缓冲型的channel if c.dataqsiz == 0 { if raceenabled { racesync(c, sg) } // 未忽略接受的数据 if ep != nil { // 直接拷贝数据,从sender goroutine -> receiver goroutine recvDirect(c.elemtype, sg, ep) } } else { // 缓冲型的channel,buf已满,将循环数组buf队首的元素拷贝到接收数据的地址,发送者数据入队,此时recvx=sendx // 找到接收游标 qp := chanbuf(c, c.recvx) ... // 将接收游标处的数据拷贝给接收者 if ep != nil { typedmemmove(c.elemtype, ep, qp) } // 将发送者数据拷贝到buf typedmemmove(c.elemtype, qp, sg.elem) // 更新游标值 c.recvx++ if c.recvx == c.dataqsiz { c.recvx = 0 } c.sendx = c.recvx // c.sendx = (c.sendx+1) % c.dataqsiz } sg.elem = nil gp := sg.g // 解锁 unlockf() gp.param = unsafe.Pointer(sg) sg.success = true if sg.releasetime != 0 { sg.releasetime = cputicks() } // 唤醒发送的goroutine,需要等待调度器光临 goready(gp, skip+1) }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如果是非缓冲型的,直接从发送者的栈拷贝到接收者的栈。
func recvDirect(t *_type, sg *sudog, dst unsafe.Pointer) { src := sg.elem typeBitsBulkBarrier(t, uintptr(dst), uintptr(src), t.size) memmove(dst, src, t.size) }1
2
3
4
5否则,属于缓冲型
channel和buf满的情况,此时发送右边和接收游标重合,需要先找到接收游标。// chanbuf(c, i) is pointer to the i'th slot in the buffer. func chanbuf(c *hchan, i uint) unsafe.Pointer { return add(c.buf, uintptr(i)*uintptr(c.elemsize)) }1
2
3
4将此处的元素拷贝到接收地址,将发送者待发送的数据拷贝到接收游标处,这样就完成了接收数据和发送数据的操作。接着,分别将发送游标和接收游标向前进一,如果发生
环绕,再从0开始。最后,会取出sudog里的goroutine,调用goready将其状态改成runnable,将发送者唤醒,等待调度器调度。func goroutineA(a <-chan int) { val := <- a fmt.Println("G1 received data: ", val) return } func goroutineB(b <-chan int) { val := <- b fmt.Println("G2 received data: ", val) return } func main() { ch := make(chan int) go goroutineA(ch) go goroutineB(ch) ch <- 3 time.Sleep(time.Second) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19创建一个无缓冲的
channel,接着启动两个goroutine,并将前面创建的channel传递进去,由主协程向channel发送数据。开始的时候,channel内什么都没有。
两个子协程被创建后,各自执行接收操作,通过前面的分析可知,
G1和G2都会被阻塞在接收操作,G1和G2会挂在channel的recq队列中,形成一个双向循环列表,sendq没有被阻塞的协程。xxxxxxxxxx func string2bytes(s string) []byte { return ([]byte)(unsafe.Pointer(&s))}func bytes2string(b []byte) string{ return *(*string)(unsafe.Pointer(&b))}go
# 2.4.chan关闭
chan关闭时,会执行函数closechan,用于修改chan状态并关闭。func closechan(c *hchan) { // 关闭一个nil chan-->panic if c == nil { panic(plainError("close of nil channel")) } // 上锁 lock(&c.lock) // 如果channel已经关闭 if c.closed != 0 { unlock(&c.lock) // panic panic(plainError("close of closed channel")) } ... // 修改关闭状态 c.closed = 1 var glist gList // 释放所有等待接收队列里的sudog for { // 从接收队列里出队一个sudog sg := c.recvq.dequeue() // 出队完毕,跳出循环 if sg == nil { break } // 如果elem不为空,说明此接收者未忽略接收数据,赋予相应类型零值 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()) } // push到链表 glist.push(gp) } // 将channel等待发送队列里的sudog释放,如果存在这些协程会panic for { // 从发送队列里取出一个sudog sg := c.sendq.dequeue() if sg == nil { break } // 发送者会panic 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()) } // push到链表 glist.push(gp) } // 解锁 unlock(&c.lock) // 遍历链表 for !glist.empty() { // 取最后一个 gp := glist.pop() // 向前走一步,下一个唤醒的g gp.schedlink = 0 // 唤醒相应的协程 goready(gp, 3) } }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
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83关闭
channel时,recvq和sendq中分别保存了阻塞的发送者和接收者,对于等待接收者而言,会收到一个相应类型的零值;对于等待发送者,会直接panic。因此,对于channel未知接收者情况下,不能贸然关闭channel。close函数先上一把大锁,接着把所有挂在这个channel上的sender和recever全部连成一个sudog链表,然后再解锁,最后将所有sudog全部唤醒,进行扫尾工作。操作 nil channel closed channel not nil, not closed channel close panic panic 正常关闭 读 <- ch 阻塞 读到对应类型的零值 阻塞或正常读取数据。缓冲型 channel 为空或非缓冲型 channel 没有等待发送者时会阻塞 写 ch <- 阻塞 panic 阻塞或正常写入数据。非缓冲型 channel 没有等待接收者或缓冲型 channel buf 满时会被阻塞
# 2.5.chan优雅关闭
channel不能重复关闭,更不建议从接受侧关闭,避免存在sender向channel发送数据时抛出panic,更不建议存在多个sender的情况下关闭channel。一般情况下,channel优雅关闭会借助另一个channel传递关闭信号实现,receiver通过信号channel下达关闭数据channel指令,senders监听到关闭信号后,停止发送数据。func main() { rand.Seed(time.Now().UnixNano()) const Max = 100000 const NumSenders = 1000 dataCh := make(chan int, 100) stopCh := make(chan struct{}) // senders for i := 0; i < NumSenders; i++ { go func() { for { select { case <- stopCh: return case dataCh <- rand.Intn(Max): } } }() } // the receiver go func() { for value := range dataCh { if value == Max-1 { fmt.Println("send stop signal to senders.") close(stopCh) return } fmt.Println(value) } }() select { case <- time.After(time.Hour): } }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这里的
stopCh就是信号channel,它本身只有一个sender,因此可以直接关闭它。senders收到关闭信号后,select分支case <- stopCh命中,退出函数不再发送数据。虽然上述程序没有显示关闭数据channel,但一个channel没有任何goroutine引用时,GC也会对其回收。
# 2.6.chan元素收发本质
channel的发送和接收操作本质上都是值的拷贝,无论从sender goroutine的栈到chan buf,还是从chan buf到receiver goroutine,或者直接从sender goroutine到receiver goroutine。
chan可能会引发goroutine泄漏,泄漏的原因是goroutine操作channel后,处于发送或接收阻塞状态,而chan处于满或空的状态,一直得不到改变。同时,垃圾回收器不会回收此类资源,进而导致goroutine会一直处于等待队列。
# 2.7.chan应用
- 并发场景下,
happened-before关系限制非常重要,编译器、CPU一般会进行各种优化,包括编译器重排、内存重排等,chan的应用遵守以下原则:- 第
n个send一定happened before第n个receiver finished,无论缓冲型还是非缓冲型 - 容量为
m的缓冲型channel,第n个receiver一定happened before第n+m个sender finished - 非缓冲型的
channel,第n个receiver一定happened before第n个send finished channel close一定happened before receiver得到通知
- 第
channel和goroutine的结合可以实现强大的并发编程能力,一般能与select、cancel、timer等结合实现信号停止、任务定时、解耦和并发控制能力。