取名好难啊啊头像
关注

Java 转 go 学习 - 并发编程(1)

本系列文章:



1. 协程

1.1 基本概念

Go 最大的特点之一就是协程,协程可以看作轻量级的线程,下面来看下进程、线程、协程之间的关系。

首先是进程,进程是 操作系统分配资源(内存、CPU 时间片、文件句柄等)的最小单位,是程序的一次运行过程(比如双击打开微信,就创建了一个微信进程)。

  • 完全独立: 每个进程都有自己的内存空间,进程之间不能直接访问,进程之间不能直接访问对方的内存,比如微信进程就不能都浏览器进程的内存。
  • 开销比较大: 创建 / 销毁 / 切换进程都需要内核参与,要拷贝大量资源比如内存映射的情况下,切换成本可能到 GB。
  • 数量有限: 一台机器允许的进程有限,通常就是几百上千。

线程就是进程内的执行单元,也是 CPU 调度和执行的 最小单位,一个进程内的线程共享这个进程的所有文件句柄、内存等资源,也就是说一个进程内可以有多个线程。

  • 资源共享: 就像上面说的,进程内的线程可以共享一个进程的堆数据、全局变量等。
  • 调度方式: 操作系统内核调度(内核态),切换的时候需要将线程里面的上下文(寄存器、程序计数器)保存起来,方便下一次切换回来的时候能够恢复现场,切换成本比较低,通常是 MB 级别的。
  • 数量上限: 能创建多少线程主要看内存大小,但是线程一般不大,像 linux 分配给线程的大小就是 8M,其他的系统可能不一样,但是基本不是很大,可以承受几千上万的线程数,比进程多。

而对于 协程 就更加轻量了,也叫微线程,协程最本质的特点就是运行在 用户态,而不是 内核态。如果说线程是操作系统调度的基本单位,那么协程就是 程序(运行时库)调度的最小单位。

  • 协程的调度: 由应用程序内部的 Runtime 或库负责 (比如 Go 的 runtime 自己调度 goroutine),切换完全在 用户态完成,只需要保存少量的寄存器状态和栈指针,不需要切换内核态,所以协程切换的开销很小(纳秒级别),比线程切换快得多了,成本是 KB 级别的。
  • 依附线程: 一个线程可以承载成千上万个协程,协程的运行必须依赖线程**(比如 1 个线程可以跑 1000 个 goroutine**。
  • 开销很小: 上面也说了,协程就是轻量级的线程,对于 Go 来说,创建一个 goroutine 仅仅需要 几 KB 空间, 一台机器能轻松运行 百万级 的 goroutine。
  • 非抢占式(Go 除外):大部分协程是协作式调度,也就是协程自己让出 CPU 才会切换 (例如遇到 await、yield 或等待 I/O 时),但是 Go 的 goroutine 是抢占式调度,runtime 会强制让出,避免某个 goroutine 霸占线程。对于协作式调度,主要问题是要避免协程陷入死循环导致整个线程阻塞,其他协程运行不了,Go 就是结合了抢占式调度可以在特定情况下强制切换协程,但是基础模型还是协作式的。

下面生成一个表格,总结下线程和协程的对比。

特性线程 (Thread)协程 (Coroutine)
调度者操作系统内核 (Kernel)用户态程序/运行时 (User Space)
切换开销高 (需内核态切换,保存完整上下文)极低 (仅用户态,保存少量寄存器)
内存占用大 (默认栈通常 1MB - 8MB)小 (初始栈通常几 KB,可动态伸缩)
并发数量有限 (受限于内存,通常几千个)极高 (单机可达百万/千万级)
编程模型抢占式 (OS随时可能中断)协作式 (需代码显式挂起)
多核利用天然支持 (OS自动分配核)需配合线程池 (将协程映射到多个线程)
锁竞争需要复杂的锁机制 (Mutex, Semaphore)单线程内无竞争,跨线程需特殊处理

1.2 Go 调度

首先有一个参数 GOMAXPROCS,这个参数是 Go 语言运行时(Runtime)中的一个核心参数,决定了同一时刻可以有多少个操作系统线程(M)同时执行 Go 代码。简单来说就是控制了 Go 程序能利用多少个 CPU 核心进行并行计算,比如下面的 CMP 模型中 GOMAXPROCS = N 就意味着 P 的最大值是 N。

Go 语言通过 GMP 模型去调度协程:

  • G:用户态的协程 Goroutine。
  • M:操作系统内核线程。
  • P:逻辑处理器,负责把 G 分配给 M 执行。

如果 GOMAXPROCS = 1:所有 Goroutine 都在同一个操作系统线程上并发运行(通过时间片切换),无法利用多核并行,只能单核运行。

如果 GOMAXPROCS = N:Go 运行时会创建 N 个 P,最多允许 N 个 M 同时绑定 P 并执行 G,从而实现在 N 个 CPU 核心上的并行运行。简单来说就是 Go 程序中最多有 N 个 goroutine 能同时在不同 CPU 核心上面执行,其他的 goroutine 就需要等待,由 Go 运行的时候调度切换。

通常情况下有 n 个核心,这个参数会设置成 n-1 来获取最好的性能,同时也需要确保 协程数 > 1 + GOMAXPROCS > 1,如果在某一个时间内只有一个协程在执行,最好就不要设置 GOMAXPROCS。

这个参数的值在不同版本的 Go 有不同的变化。

  • Go 1.4 及之前: 默认值为 1,也就是说如果不手动设置,Go 程序永远只会在一个核上跑,哪怕你有 64 核的服务器,开发者必须手动调用 runtime.GOMAXPROCS(runtime.NumCPU())。
  • Go 1.5 及之后: 默认值改成 runtime.NumCPU(),也就是机器可用的逻辑 CPU 核心数,Go 程序启动时会自动检测硬件配置,利用多核 CPU 的性能。
  • Go 1.25+ 及之后: 容器化环境如 Docker, Kubernetes 中,由于容器被限制了 CPU 使用量,比如限制成 0.5 核或者 2 核,但是 Go 运行时候读取的还是宿主机的总核心数,比如 16 核。后果就是 GOMAXPROCS 默认为 16,导致 16 个线程去抢 0.5 核个 CPU 的 时间片,造成上下文切换太频繁,性能下降,反而比单线程还慢,所以 1.25+ 之后 Go 运行时原生支持 Container-aware GOMAXPROCS,如果检测到运行在容器中,会直接读取 Cgroup 的 CPU 限制来设置默认值,不再需要第三方库如 uber-go/automaxprocs 干预。

下面是一些应用场景,不过默认值(等于 CPU 核心数)在 99% 的场景下都是最优的。

  1. 计算密集型任务: 保持默认或略低于核心数,这类任务的特点是 goroutine 几乎不会阻塞,会一直占用 CPU。

    • 如果 GOMAXPROCS 等于 CPU 核数,就i是每个核心跑一个 CPU,这种情况下没有内核线程切换开销,性能最优。
    • 如果 GOMAXPROCS 大于 CPU 核数,OS 会对 M 线程做内核态切换,反而增加开销。
  2. IO 密集型任务:默认值即可(无需调整),这类任务(比如网络请求、文件读写、数据库操作)特点是 goroutine 大部分时间都在阻塞(等待 IO 响应)。理解一个概念就是 GOMAXPROCS 管的是「CPU 执行」,不管 「IO 阻塞」,也就是说大部分时间 goroutine 需要消耗 CPU 的是发起网络请求和处理响应结果,中间的等待响应是让出 CPU 的,默认值已经能够处理消耗 CPU 的任务了。

    • 当 goroutine 阻塞的时候,P 回解绑当前 M,绑定新的 M 的执行队列中的下一个 G。
    • GOMAXPROCS=8 能支持上万的并发 goroutine,因为大部分时间 G 都在阻塞,不会占用 CPU。
    • 这种情况下不需要增大 GOMAXPROCS,增大反而会创建更多 M,比如 8 和 CPU 却将这个参数设置成 16,这样就会导致 16 个 P 会绑定 16 个线程,系统需要在 16 个线程之间轮流调度资源,抢占 8 个 CPU 核心,其中必然涉及到大量线程切换,属于内核态开销,增加内核压力。
  3. 容器 / 虚拟化环境:手动设置为容器分配的核心数

    • 这种情况上面已经说过了,就不多复述,可以手动设置 GOMAXPROCS 等于容器分配的核心数。

下面我们可以通过代码看下这个参数的值。

package main

import (
	"fmt"
	"runtime"
)

func main() {
	// 打印当前 GOMAXPROCS 值
	// GOMAXPROCS = 8
	fmt.Printf("GOMAXPROCS = %d\n", runtime.GOMAXPROCS(0))
	// 打印宿主机逻辑 CPU 核心数
	// Host CPU cores = 8
	fmt.Printf("Host CPU cores = %d\n", runtime.NumCPU())
}

1.3 基本用法

上面两个小节学习了协程的一些基本概念,这个小节正式进入协程的用法,来看下代码里面如何创建执行一个协程。

协程的启动很简单,只需在函数调用前加一个关键字 go。

go functionName(args...)

下面做一个简单的示例,我们来并发打印。

func test2() {
	for i := 0; i < 3; i++ {
		go sayHello(i)
	}
	// 阻塞等待 goroutine 打印
	time.Sleep(time.Second)
}

func sayHello(id int) {
	for i := 0; i < 3; i++ {
		fmt.Printf("Goroutine %d: Hello %d\n", id, i)
		time.Sleep(100 * time.Millisecond) // 耗时 100ms
	}
}

输出如下:

Goroutine 2: Hello 0
Goroutine 1: Hello 0
Goroutine 0: Hello 0
Goroutine 0: Hello 1
Goroutine 2: Hello 1
Goroutine 1: Hello 1
Goroutine 1: Hello 2
Goroutine 2: Hello 2
Goroutine 0: Hello 2

2. Channel 通信

2.1 基本概念

Channel 就是连接 Goroutine 的管道,协程之间如果想要互相通信,可以用共享变量,但是这样就要确保共享变量的线程安全问题。

Go 提供了 channel 类型,这也是一种类型,可以通过 channel 来进行数据的发送和读取,就类似 管道,通道通信首先能确保同步性,发送和接收默认都是阻塞的,不需要加锁。且同一时间内只有一个协程可以访问到数据,不会出现数据竞争。总结下 channel 的特点:

  • 类型安全: 一个管道的数据类型是固定的,chan int 只能传 int,chan string 只能传 string。
  • 自动同步: 发送和接收默认是阻塞的,同一时间只能由一个协程可以读写数据,天然实现同步,不需要加锁。

var identifier chan datatype 可以用来声明一种管道类型,如果想要创建 channel 可以用 make 关键字。

  • ch := make(chan int):创建无缓冲的通道。
  • ch := make(chan int, 5):创建有缓冲的通道。

缓冲意思就是管道里面能存多少数据。

  • 无缓冲 channel: 发送数据会阻塞,直到有 goroutine 接收,接收数据也会阻塞,直到有 goroutine 发送。
  • 有缓冲 channel: 只有缓冲满时发送才阻塞,缓冲空时接收才阻塞。

创建出缓冲之后可以通过 close() 关闭 channel,关闭之后就没办法再往里面发送数据,但是一样可以接收剩余的数据。

通道之间的通信操作符是 <-,比较直观,信息就按照箭头的方向流动。

  • ch <- int1 表示:将变量 int1 发送到通道 ch 中。
  • int2 := <- ch 表示:从通道 ch 读取数据到 int2。

2.2 基础例子

有了这个下面我们可以来看下例子,分别是有缓冲和无缓冲。

func test3() {
	// 1. 创建无缓冲 channel
	ch := make(chan int)

	// 2. 启动 goroutine 发送数据
	go func() {
		fmt.Println("goroutine: 准备发送数据")
		ch <- 100 // 没有缓冲阻塞等待读取
		fmt.Println("goroutine: 数据发送完成")
	}()

	time.Sleep(time.Second)
	// 3. goroutine 接收数据
	fmt.Println("main: 准备接收数据")
	num := <-ch // 从 channel 接收数据
	fmt.Println("main: 接收到数据 =", num)

	// 4. 关闭 channel
	close(ch)
	
	// goroutine: 准备发送数据
	// main: 准备接收数据
	// main: 接收到数据 = 100
	// goroutine: 数据发送完成
}

上面是无缓冲的例子,发送数据之后需要一直阻塞直到将 channel 里面的数据读掉,下面再来看下有缓冲的例子。

func test3() {
	// 创建缓冲大小为 2 的 channel
	ch := make(chan string, 2)

	// 发送 2 条数据: ,不会阻塞
	ch <- "hello"
	ch <- "channel"
	fmt.Println("发送 2 条数据完成, 缓冲已满")

	// ch <- "go" // 发送第三条会阻塞
	time.Sleep(time.Second)
	// 接收数据
	fmt.Println("接收 1 条数据:", <-ch)
	fmt.Println("接收 1 条数据:", <-ch)

	// 发送 2 条数据完成, 缓冲已满
	// 接收 1 条数据: hello
	// 接收 1 条数据: channel
	close(ch)
}

可以看到我们上面创建了大小为 2 的缓冲区,然后往里面发送数据,发送了 2 条之后缓冲区就满了,如果这时候再发送一条就会阻塞,上面这种情况就不会。

当然除了上面这种,我们也可以在 close 关闭之后遍历 channel 里面的数据来处理,注意如果要用 range 遍历一定要 close,否则会一直阻塞。

func test5() {
	ch := make(chan int, 3)
	// 发送数据
	ch <- 1
	ch <- 2
	ch <- 3
	close(ch) // 必须关闭, 否则 range 会一直阻塞等待新数据

	// 遍历 channel: 自动读取所有数据直到 channel 关闭
	for num := range ch {
		fmt.Println("遍历 channel ", num)
	}
	
	// 遍历 channel  1
	// 遍历 channel  2
	// 遍历 channel  3
}

2.3 信号量

通过通道的特性,我们可以启动一个简单的信号量,缓冲大小 = 信号量的最大并发数,获取信号量就是从 channel 里面获取数据,计数器 - 1,释放信号量就是向 channel 里面放数据,计数器 - 1。

下面我们定义一个信号量结构体,然后结构体里面的属性就是 chan,类型是 struct{},空结构体不占内存,省空间。

type Semaphore struct {
	// 空结构体, 占用 0, 节省资源
	ch chan struct{}
}

然后定义信号量的初始化方法,传入信号量数量大小。

func NewSemaphore(size int) *Semaphore {
	if size <= 0 {
		// 最少是 1
		size = 1
	}
	return &Semaphore{make(chan struct{}, size)}
}

下面接着定义获取信号量和释放信号量的方法。

// 获取信号量
func (s *Semaphore) Acquire() {
	// struct{}{} 创建空结构体
	s.ch <- struct{}{}
}

// 释放信号量
func (s *Semaphore) Release() {
	<-s.ch
}

上面获取信号量的时候通过 struct{}{} 去创建一个空的结构体。

最后测试如下。


// 获取格式化时间:YYYY-MM-DD HH:mm:ss:SSS
func getCurrentTime() string {
	return time.Now().Format("2006-01-02 15:04:05.000")
}

func test21() {
	const (
		maxConcurrent = 3 // 最大并发数(信号量大小)
		totalTask     = 8 // 总任务数
	)

	// 1. 创建信号量, 最大 3 个并发
	sem := NewSemaphore(maxConcurrent)
	var wg sync.WaitGroup
	wg.Add(totalTask)

	// 2. 启动 8 个 goroutine, 但是限制同时只能有 3 个在执行
	for i := 1; i <= totalTask; i++ {
		taskID := i
		go func() {
			defer wg.Done()

			// 获取信号量
			sem.Acquire()
			// 确保任务执行完成之后再释放信号量
			defer sem.Release()

			// 模拟任务执行
			fmt.Printf("[%s] 任务%d 开始执行(当前并发数:%d)\n", getCurrentTime(), taskID, len(sem.ch))
			time.Sleep(1 * time.Second) // 模拟耗时操作
			fmt.Printf("[%s] 任务%d 执行完成\n", getCurrentTime(), taskID)
		}()
	}

	wg.Wait()
	fmt.Printf("[%s] 所有任务执行完毕\n", getCurrentTime())
}

输出结果如下:

[2026-03-23 22:14:16.627] 任务1 开始执行(当前并发数:3)
[2026-03-23 22:14:16.628] 任务3 开始执行(当前并发数:3)
[2026-03-23 22:14:16.628] 任务8 开始执行(当前并发数:3)
[2026-03-23 22:14:17.645] 任务8 执行完成
[2026-03-23 22:14:17.645] 任务3 执行完成
[2026-03-23 22:14:17.645] 任务5 开始执行(当前并发数:3)
[2026-03-23 22:14:17.645] 任务4 开始执行(当前并发数:3)
[2026-03-23 22:14:17.645] 任务1 执行完成
[2026-03-23 22:14:17.645] 任务6 开始执行(当前并发数:3)
[2026-03-23 22:14:18.645] 任务4 执行完成
[2026-03-23 22:14:18.645] 任务6 执行完成
[2026-03-23 22:14:18.645] 任务5 执行完成
[2026-03-23 22:14:18.645] 任务7 开始执行(当前并发数:3)
[2026-03-23 22:14:18.645] 任务2 开始执行(当前并发数:2)
[2026-03-23 22:14:19.646] 任务7 执行完成
[2026-03-23 22:14:19.646] 任务2 执行完成
[2026-03-23 22:14:19.646] 所有任务执行完毕

2.4 通道方向

通道方向可以表示只接受或者只发送:

  • var send_only chan<- int:只发送
  • var recv_only <-chan int:只接收

那对于只接收的通道是没办法关闭的,因为关闭通道是发送者用来表示不再给通道发送消息了,所以对于只接收通道来说关闭时没有意义的。

首先有一点要明确的,不能直接创建有方向的,默认 make 创建的都是双向通道,但是可以通过类型转换做到有方向的通道。

通道类型声明语法允许操作禁止操作
双向通道(默认)ch chan T发送 + 接收无
只发送通道ch chan<- T仅发送接收、关闭(注)
只接收通道ch <-chan T仅接收发送

下面我们写一个例子,先创建双向通道,然后转成只写和只读通道,接下来分别从只写通道写入数据,然后从只读通道读出。

func test31() {
	// 1. 创建双向通道, 有没有缓冲都行
	biCh := make(chan string, 3)

	// 2. 类型转换
	var sendCh chan<- string = biCh
	var readCh <-chan string = biCh

	// 3. 验证
	sendCh <- "hello"
	sendCh <- "world"
	// readCh <- "hello" 报错 Invalid operation: readCh <- "hello" (send to the receive-only type <-chan string)

	// 从只读通道中读出数据, 注意 sendCh 和 readCh 相当于 biCh 这个通道的两边, 不是两个通道
	fmt.Println(<-readCh)
	fmt.Println(<-readCh)
	
	// hello
	// world

	// close(readCh) 报错 Must be a bidirectional or send-only channel
	close(sendCh)
}

2.5 检测通道是否关闭

在 Go 中,没有直接查询通道状态的函数,但是可以通过接收通道数据时第二个返回值来判断,上面几个小节的例子中用的都是一个参数,其实还有一个参数没用。

v, ok := <-ch
  • ok = true:通道没有关闭,v 是从通道接受到的有效数据。
  • ok = false:通道已经关闭,且通道内没有剩余的数据,此时的 v 是通道类型的零值。

下面来看下例子:

func test32() {
	// 1. 创建双向通道并发送数据
	ch := make(chan int, 2)
	ch <- 10
	ch <- 20
	close(ch) // 关闭通道

	// 2. 第一次接收, 通道里面还有数据
	v1, ok1 := <-ch
	fmt.Printf("接收值:%d,通道是否关闭:%t(ok=%t)\n", v1, !ok1, ok1)
	// 接收值:10,通道是否关闭:false(ok=true)

	// 3. 第二次接收, 通道里面还有数据
	v2, ok2 := <-ch
	fmt.Printf("接收值:%d,通道是否关闭:%t(ok=%t)\n", v2, !ok2, ok2)
	// 接收值:20,通道是否关闭:false(ok=true)

	// 4. 第三次接收, 通道已关闭且无数据
	v3, ok3 := <-ch
	fmt.Printf("接收值:%d,通道是否关闭:%t(ok=%t)\n", v3, !ok3, ok3)
	// 接收值:0,通道是否关闭:true(ok=false)

	// 5. 对未关闭的通道检测(对比)
	ch2 := make(chan string)
	go func() {
		time.Sleep(100 * time.Millisecond)
		ch2 <- "hello" // 未关闭通道发送数据
	}()
	v4, ok4 := <-ch2
	fmt.Printf("未关闭通道接收值:%s,ok=%t\n", v4, ok4)
	// 未关闭通道接收值:hello,ok=true

}

for range 遍历通道时,会自动退出循环,不需要手动判断。

func test33() {
	// 创建双向通道并发送数据
	ch := make(chan int, 2)
	ch <- 10
	ch <- 20
	close(ch) // 关闭通道

	for i := range ch {
		fmt.Println(i)
	}
	// 10
	// 20
}

最后再来说下一些注意事项,首先就是 禁止重复关闭通道,如果重复关闭通道会报错,这种情况下可以使用 sync.Once 确保通道只关闭一次。

func test34() {
	ch := make(chan int, 2)
	ch <- 10
	ch <- 20
	close(ch) // 关闭通道
	close(ch) // 关闭通道: panic: close of closed channel
}

然后就是不要向已经关闭的通道发送消息,同样会 panic。

func test34() {
	ch := make(chan int, 2)
	ch <- 10
	ch <- 20
	close(ch) // 关闭通道
	ch <- 20 // panic: send on closed channel
}

最后就是没有经过初始化的通道 var ch chan int 接收发送都会永久阻塞,v, ok := <-ch 也会阻塞,没办法检查状态,需要用 ch = make(chan int) 来初始化才行。


2.6 Select 语句 - 多路复用

2.6.1 基础概念

如果我们需要同时监听 多个 Channel 或者 处理超时,就需要用到 select,select 类似 switch,但是专门用来操作 Channel 的,一旦某个通道可执行 (发送 / 接收就绪),就执行对应的 case 逻辑;如果多个通道同时就绪,随机选择一个执行(没有优先级)。

核心特点如下:

  1. Select 只用来操作 Channel: case 只能是通道的发送 ch <- v 或者接收 <-ch 操作,不是普通的条件比如 i > 10 这种。
  2. Select 处理通道没有优先级:如果多个 case 同时就绪,那么 Go 会随机选择一个通道来执行。
  3. 阻塞 / 非阻塞可控:没有就绪 case 并且没有 default 的时候,select 会阻塞,如果有 default 就会执行 default,这种情况就是非阻塞。
  4. 可配合循环: 单次 select 只执行一个 case,如果要处理多个 channel 需要循环执行,所以一般都是 for-select 配合使用。

基本语法如下:

select {
case <-ch1: // 通道ch1接收数据就绪
    // 处理逻辑
case ch2 <- v: // 通道ch2发送数据就绪
    // 处理逻辑
case v, ok := <-ch3: // 接收+检测通道是否关闭
    // 处理逻辑
default: // 所有case都未就绪时执行(非阻塞)
    // 兜底逻辑
}

这里要注意如果是一个空的 select,也就是没有 case,这种情况下会永久阻塞当前的 goroutine,等于 for{}。

下面来看一个例子,监听多个通道。

func test35() {
	ch1 := make(chan string)
	ch2 := make(chan int)

	// goroutine1 向 ch1 发送数据
	go func() {
		// 1s
		time.Sleep(1 * time.Second)
		ch1 <- "hello from ch1"
	}()

	// goroutine2 向 ch2 发送数据
	go func() {
		// 2s
		time.Sleep(2 * time.Second)
		ch2 <- 100
	}()

	fmt.Println("开始监听 ch1 和 ch2...")
	select {
	case v := <-ch1:
		fmt.Println("收到 ch1 的数据:", v)
	case v := <-ch2:
		fmt.Println("收到 ch2 的数据:", v)
	}
	fmt.Println("select 执行完成")
}

func main() {
	test35()
	// 开始监听 ch1 和 ch2...
	// 收到 ch1 的数据: hello from ch1
	// select 执行完成
}

可以看到因为 ch1 睡眠 1s,所以比 ch2 更早就绪,因此 select 会处理 ch1 的 case。

下面我们来看下往一个已经满了的通道发送消息的场景,使用 select 来进行兜底。

func test36() {
	ch := make(chan int, 1)
	ch <- 5

	// 非阻塞发送, 通道已满, 执行 default 分支
	select {
	case ch <- 5:
		fmt.Println("往 ch 写入 5")
	default:
		fmt.Println("通道数据已满, 发送失败")
	}
	
	// 通道数据已满, 发送失败
}

2.6.2 超时处理

接下来就是重点超时处理了,如果一个通道长时间没有发送数据,总不可能一直阻塞,这时候就可以结合 time.After 加上超时处理。

func test37() {
	ch := make(chan string)

	go func() {
		// 模拟永远不往 ch 里面发送数据
		time.Sleep(10 * time.Second)
	}()

	select {
	case v := <-ch:
		fmt.Println("收到 ch 消息:", v)
	case <-time.After(1 * time.Second):
		fmt.Println("1s 时间内都没有收到 ch 消息, 退出")
	}
}

func main() {
	test37()
	// 1s 时间内都没有收到 ch 消息, 退出
}

上面的例子中,当 1s 后还没有收到消息,第二个 case 就会收到消息打印日志退出,并不会长时间阻塞。


2.6.3 持续监听 + 退出信号

select 一次只能处理一个 case,因此我们需要结合 for 来实现 持续监听多个通道,并且通过退出通道来控制停止。

func test38() {
	workCh := make(chan int)
	quitCh := make(chan struct{}) // 退出信号, 使用空结构体省内存

	// 工作 goroutine 不断往里面写数据
	go func() {
		for i := 1; ; i++ {
			workCh <- i
			time.Sleep(500 * time.Millisecond)
		}
	}()

	// 主 goroutine 持续监听
	fmt.Printf("[%s] 开始工作, 5s 后退出\n", util.Now())
	go func() {
		time.Sleep(5 * time.Second)
		quitCh <- struct{}{}
	}()

	// for-select 持续监听
	for {
		select {
		case num := <-workCh:
			fmt.Printf("[%s] 收到工作 goroutine 写入的数据: %d\n", util.Now(), num)
		case <-quitCh:
			fmt.Printf("[%s] 收到退出信号, 退出\n", util.Now())
			return
		}
	}
	
	// [2026-03-24 23:41:57] 开始工作, 5s 后退出
	// [2026-03-24 23:41:57] 收到工作 goroutine 写入的数据: 1
	// [2026-03-24 23:41:58] 收到工作 goroutine 写入的数据: 2
	// [2026-03-24 23:41:58] 收到工作 goroutine 写入的数据: 3
	// [2026-03-24 23:41:59] 收到工作 goroutine 写入的数据: 4
	// [2026-03-24 23:41:59] 收到工作 goroutine 写入的数据: 5
	// [2026-03-24 23:42:00] 收到工作 goroutine 写入的数据: 6
	// [2026-03-24 23:42:00] 收到工作 goroutine 写入的数据: 7
	// [2026-03-24 23:42:01] 收到工作 goroutine 写入的数据: 8
	// [2026-03-24 23:42:01] 收到工作 goroutine 写入的数据: 9
	// [2026-03-24 23:42:02] 收到工作 goroutine 写入的数据: 10
	// [2026-03-24 23:42:02] 收到退出信号, 退出
}

2.6.4 检测通道关闭(结合 ok 值)

我们可以在 select 中结合通道读取返回的 ok 值来判断一条通道是不是关闭了,下面来看一个例子。

func test39() {
	ch := make(chan int)
	var wg sync.WaitGroup

	// 生产者: 发送 3 个数据之后关闭通道
	go func() {
		for i := 0; i < 3; i++ {
			ch <- i
		}
		close(ch)
		fmt.Println("[生产组] 数据发送已完成, 关闭通道")
	}()

	// 消费者 select + ok 检测通道关闭
	wg.Add(1)
	go func() {
		defer wg.Done()
		for {
			select {
			case num, ok := <-ch:
				if !ok {
					fmt.Println("[消费者] 通道已关闭, 退出程序")
					return
				} else {
					fmt.Println("[消费者] 收到数据:", num)
				}
			}
		}
	}()

	wg.Wait()
	
	// [消费者] 收到数据: 0
	// [消费者] 收到数据: 1
	// [消费者] 收到数据: 2
	// [生产组] 数据发送已完成, 关闭通道
	// [消费者] 通道已关闭, 退出程序
}

3. 小结

这篇文章中我们学习了 go 语言的并发编程基础,包括一些编程的注意事项,下一篇文章继续学习并发编程相关的内容。





如有错误,欢迎指出!!!!

转载自 CSDN-专业IT技术社区

原文链接:https://blog.csdn.net/laohuangaa/article/details/159354746

文章来源转载

评论

赞0

评论列表

微信小程序
QQ小程序

关于作者

点赞数:0
关注数:0
粉丝:0
文章:0
关注标签:0
加入于:--