Go如何实现“等待 100 个 goroutine 完成,但最多等待 3 秒“

context取消

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
package main

import (
"context"
"fmt"
"sync"
"time"
)

func main() {
var wg sync.WaitGroup
// 创建一个 3 秒后自动取消的 Context
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel() // 确保在 main 退出时释放资源

for i := 0; i < 100; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
// 模拟工作:随机耗时 1-5 秒
workTime := time.Duration(id%4+1) * time.Second

select {
case <-time.After(workTime):
// 正常完成工作
fmt.Printf("Goroutine %d 完成\n", id)
case <-ctx.Done():
// 收到超时信号,立即退出
fmt.Printf("Goroutine %d 超时退出\n", id)
return
}
}(i)
}

// 等待所有 goroutine 退出(无论是完成还是超时)
wg.Wait()
fmt.Println("所有 goroutine 已终止,程序退出")
}

time.After 和 done channel(主控端强杀)

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
package main

import (
"fmt"
"sync"
"time"
)

func main() {
var wg sync.WaitGroup
done := make(chan struct{}) // 用于通知所有 goroutine 退出

// 启动 100 个 goroutine
for i := 0; i < 100; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
// 模拟工作,循环中必须检查 done 信号
for {
select {
case <-done:
fmt.Printf("Goroutine %d 收到退出信号\n", id)
return
default:
// 模拟执行一小段任务
time.Sleep(100 * time.Millisecond)
fmt.Printf("Goroutine %d 工作中...\n", id)
}
}
}(i)
}

// 主协程等待 3 秒
select {
case <-time.After(3 * time.Second):
fmt.Println("3秒超时,发送退出信号")
close(done) // 广播退出信号给所有 goroutine
}

// 等待所有 goroutine 响应退出信号并清理
wg.Wait()
fmt.Println("所有 goroutine 已处理完毕")
}

如果现在有一个总goroutine,在内部会启动多个goroutine,这个总goroutine要怎样结束内部启动的这些goroutine?

context.WithCancel

这是 Go 官方推荐的父子任务管理方式。父 goroutine 持有 cancel 函数,当需要退出时,调用它,所有子 goroutine 监听 ctx.Done() 通道。

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
package main

import (
"context"
"fmt"
"sync"
"time"
)

func main() {
var wg sync.WaitGroup
// 1. 创建可取消的 Context
ctx, cancel := context.WithCancel(context.Background())

// 2. 启动总 goroutine(父)
wg.Add(1)
go func() {
defer wg.Done()
fmt.Println("总 goroutine 启动,开始派生子任务...")

// 启动多个子 goroutine
for i := 0; i < 3; i++ {
wg.Add(1)
go func(id int) {
defer wg.Done()
for {
select {
case <-ctx.Done(): // 3. 监听取消信号
fmt.Printf("子 goroutine %d 收到退出信号,正在清理并退出\n", id)
return
default:
// 模拟工作
fmt.Printf("子 goroutine %d 工作中...\n", id)
time.Sleep(500 * time.Millisecond)
}
}
}(i)
}

// 模拟总 goroutine 运行一段时间后结束
time.Sleep(2 * time.Second)
fmt.Println("总 goroutine 任务完成,准备通知所有子 goroutine 退出...")
// 4. 调用取消函数
cancel()
}()

// 等待所有 goroutine 优雅退出
wg.Wait()
fmt.Println("所有 goroutine 已安全退出,主程序结束")
}

go里面有一个参数并不确定类型,应该怎样去定义它?

空接口 interface{} (any)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
func Process(data interface{}) {
// 使用类型断言判断具体类型
switch v := data.(type) {
case string:
fmt.Println("字符串:", v)
case int:
fmt.Println("整数:", v)
case bool:
fmt.Println("布尔:", v)
default:
fmt.Printf("未知类型: %T\n", v)
}
}

// Go 1.18+ 可以使用 any 别名
func Process2(data any) {
// 同样的处理方式
}

泛型 (Go 1.18+)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
// 约束为可比较类型
func Process[T comparable](data T) {
fmt.Printf("值: %v, 类型: %T\n", data, data)
}

// 约束为特定接口
func ProcessNumber[T int | float64](data T) T {
return data * 2
}

// 使用自定义约束
type Addable interface {
int | int64 | float64 | string
}

func Add[T Addable](a, b T) T {
return a + b
}

JSON.RawMessage

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
type Request struct {
Type string `json:"type"`
Payload json.RawMessage `json:"payload"` // 延迟解析
}

func Handle(req Request) {
switch req.Type {
case "user":
var user User
json.Unmarshal(req.Payload, &user)
// 处理 user
case "order":
var order Order
json.Unmarshal(req.Payload, &order)
// 处理 order
}
}

map[string]interface{}

1
2
3
4
5
6
7
8
func ProcessMap(data map[string]interface{}) {
if val, ok := data["name"].(string); ok {
fmt.Println("Name:", val)
}
if val, ok := data["age"].(float64); ok {
fmt.Println("Age:", int(val))
}
}

定义接口约束

如果参数需要实现某些方法:

1
2
3
4
5
6
7
8
9
10
type Processor interface {
Process() string
}

func Handle(p Processor) {
result := p.Process()
fmt.Println(result)
}

// 任何实现了 Process() string 的类型都可以传入

避免过度使用 interface{},会失去类型安全

泛型是编译时多态,比 interface{} 性能更好

使用类型断言时总是检查 ok,避免 panic

考虑使用 类型开关 (type switch) 处理多种情况

接口有啥用?

假设你在写一个用户管理系统,要保存用户数据。今天用 MySQL,明天可能换 Redis,后天换 MongoDB。

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
package main

import "fmt"

// 1. 定义存储接口:你要能存用户,能查用户
type UserStorage interface {
Save(name string)
Find(name string) string
}

// 2. 实现一:MySQL 版
type MySQL struct{}
func (m MySQL) Save(name string) {
fmt.Printf("MySQL 保存了用户:%s\n", name)
}
func (m MySQL) Find(name string) string {
return fmt.Sprintf("MySQL 查到了用户:%s", name)
}

// 3. 实现二:Redis 版(假装实现)
type Redis struct{}
func (r Redis) Save(name string) {
fmt.Printf("Redis 保存了用户:%s\n", name)
}
func (r Redis) Find(name string) string {
return fmt.Sprintf("Redis 查到了用户:%s", name)
}

// 4. 业务函数:只依赖接口,不依赖具体实现
func RegisterUser(storage UserStorage, name string) {
storage.Save(name)
}

func GetUser(storage UserStorage, name string) {
fmt.Println(storage.Find(name))
}

func main() {
mysql := MySQL{}
redis := Redis{}

// 用 MySQL
RegisterUser(mysql, "张三")
GetUser(mysql, "张三")

fmt.Println("-----切换数据库-----")

// 换成 Redis,业务代码一行不改!
RegisterUser(redis, "李四")
GetUser(redis, "李四")
}

在 Go 语言中,panicrecover 是内置的用于处理异常崩溃的机制。要捕捉 panic,必须使用 defer 配合 recover

1. 核心范式:defer + recover

recover() 只有在 defer 函数中调用时才有效。当 panic 发生时,当前函数会立即停止执行,转而执行所有已压栈的 defer 语句,此时 recover 才能捕获到异常。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
package main

import "fmt"

func main() {
// 使用 defer 匿名函数来捕获 panic
defer func() {
if r := recover(); r != nil {
fmt.Println("成功捕获 panic:", r)
// 在这里可以做日志记录、状态恢复等操作
}
}()

// 触发 panic
panic("发生严重错误")

// 这行代码不会执行
fmt.Println("程序继续...")
}

输出:

1
成功捕获 panic: 发生严重错误

2. 获取更详细的堆栈信息

仅仅获取 panic 的错误信息通常不够,建议使用 runtime/debug 包打印堆栈轨迹,方便排查问题。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
import (
"fmt"
"runtime/debug"
)

func SafeCall() {
defer func() {
if r := recover(); r != nil {
// 打印错误信息和完整调用栈
fmt.Printf("捕获到 Panic: %v\n堆栈信息:\n%s\n", r, debug.Stack())
}
}()

// 业务逻辑...
panic("数组越界")
}

3. 在实际项目中的最佳实践

不要main 函数的最外层捕获所有 panic(除非是守护进程),推荐的做法是在 Goroutine 的入口处关键业务逻辑边界(如 HTTP/RPC 处理函数)添加恢复逻辑。

场景一:每个 Goroutine 必须单独捕获

注意: main 函数中的 defer recover 无法捕获其他子 Goroutine 的 panic每个 Goroutine 启动时,必须单独使用 defer recover

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
func main() {
// 错误写法:这个 recover 捕获不到子协程的 panic
// defer recover()

go func() {
defer func() {
if r := recover(); r != nil {
log.Printf("Goroutine 崩溃: %v", r)
}
}()

// 业务逻辑
panic("子协程出错")
}()

// 阻塞主线程
select {}
}

场景二:Web 框架中的中间件(以 Gin 为例)

Gin 框架自带 recovery 中间件,如果你要自定义,可以这样做:

1
2
3
4
5
6
7
8
9
10
11
12
13
func RecoveryMiddleware() gin.HandlerFunc {
return func(c *gin.Context) {
defer func() {
if err := recover(); err != nil {
// 记录日志
logger.Error("接口发生 Panic", zap.Any("error", err))
// 返回 500 状态码给客户端,而不是让服务直接挂掉
c.AbortWithStatusJSON(500, gin.H{"code": 500, "msg": "服务器内部错误"})
}
}()
c.Next() // 执行后续处理函数
}
}

4. 关键注意事项(避坑指南)

注意事项 说明
仅捕获当前协程 recover 只能捕获同一 Goroutine 中发生的 panic。跨协程无法捕获。
不要滥用 不要试图用 recover 捕获所有异常当作 if err 来用。只有不可恢复的错误(如数组越界、空指针)才用 panic,业务逻辑错误请用 error 返回。
恢复后的状态 捕获 panic 后,程序会继续执行,但该函数中 panic 之后的代码不会执行,只会执行 defer 之后的后续流程。
重新抛出 如果你捕获了 panic 但无法处理(比如资源无法释放),可以再次 panic(r) 将异常向上抛出。

5. 高级技巧:区分类型

recover() 返回的是 interface{},你可以通过类型断言来处理不同类型的 Panic:

1
2
3
4
5
6
7
8
9
10
11
12
defer func() {
if r := recover(); r != nil {
switch v := r.(type) {
case string:
fmt.Println("字符串异常:", v)
case error:
fmt.Println("错误对象异常:", v.Error())
default:
fmt.Println("未知类型异常:", v)
}
}
}()

** recover 并不能真正”恢复”(recover)程序到崩溃前的状态,它只是”截获”(catch)了 panic 异常,让程序不至于彻底崩溃退出。**

为什么用grpc作为通信协议

选择gRPC,本质上是在为高性能、高要求的内部系统通信选择一条“高速公路”。它和常见的基于HTTP/1.1的RESTful API走的是完全不同的技术路线。

它的核心优势主要集中在三个方面,可以看作是“更快、更稳、更规范”。

🚀 极致性能:“快”是首要目标

gRPC追求极致的通信效率,性能远超传统的REST API。这得益于两大技术的组合:

  1. 传输层:基于HTTP/2协议:与HTTP/1.1每个请求都需要新建一个TCP连接不同,HTTP/2支持多路复用,可以在同一个TCP连接上并行处理成百上千个请求,极大地减少了连接建立的开销。同时,它的头部压缩(HPACK)和二进制分帧机制也让传输效率提升了30%-50%。学术研究也证实,在高并发场景下,gRPC的吞吐量远高于REST,响应时间也更稳定。
  2. 序列化:使用Protocol Buffers (Protobuf):与JSON这类文本格式不同,Protobuf采用二进制编码,数据体积更小(通常比JSON小3-5倍),序列化和反序列化的速度更是快上5-8倍。对于大数据量和频繁调用的场景,节省的带宽和CPU资源非常可观。

🧩 开发效率:“稳”是坚实保障

gRPC不只是快,它带来的“强契约”和“跨语言”能力,能从源头上提升团队协作和开发效率。

  • 强类型契约:gRPC通过.proto文件严格定义服务接口和数据结构。这就像一份“合同”,客户端和服务端都必须遵守,有效避免了JSON等弱类型在解析时可能出现的字段拼写错误或类型不匹配问题。
  • 自动化多语言代码生成:利用protoc编译器,只要定义好一份.proto文件,就能自动生成Java、Go、Python等10多种主流语言的客户端和服务端代码。这种“语言无关性”非常适合技术栈多元化的团队。

🔧 强大的RPC模型:满足复杂通信需求

gRPC原生支持四种通信模式,能覆盖从简单调用到复杂实时交互的各种场景:

  • 一元RPC:标准的请求-响应模式。
  • 服务端流式RPC:服务端可以持续向客户端推送数据,如实时行情。
  • 客户端流式RPC:客户端可以持续向服务端发送数据,如上传大文件。
  • 双向流式RPC:客户端和服务端可以同时、独立地发送消息流,适用于实时聊天或游戏等场景。

🤔 那么,gRPC适合所有场景吗?

当然不。由于它的设计初衷是用于“服务间通信”,也存在一些天然的局限:

  • 浏览器支持有限:浏览器JavaScript引擎无法直接支持gRPC的HTTP/2调用,需要通过额外的代理层(如gRPC-Web)转换,这会带来性能损耗和复杂性。
  • 调试相对复杂:因为数据是二进制格式的,无法像JSON那样直接在浏览器或命令行工具里查看,排查问题需要专门的工具。

📝 总结

总的来说,选择gRPC的场景非常明确:

  • 最佳选择内部的微服务通信、对性能有极致要求的系统、需要多语言协作的大型项目。
  • 需要三思:对外提供公开API(如开放平台)、需要浏览器直接调用的Web应用,以及简单的CRUD操作。在这些场景下,简单、通用的RESTful API可能是更务实的选择。

一致性哈希和普通哈希的区别

这个问题问得很及时,刚才我们聊的gRPC通常用于微服务,而一致性哈希正是解决这类分布式系统中“如何找对机器”这个核心问题的关键算法。

要理解它们的区别,最直观的方式是看一个扩容时的场景。我们用一个经典的“取模哈希”来代表普通哈希,对比看看会发生什么。

核心区别:扩容时的“灾难”与“优雅”

假设你有 3台 服务器(编号0,1,2),存了 3个 数据(A,B,C)。为了决定数据存到哪台机器,常规做法是计算 hash(数据) % 3

  • 可能的结果是:A存在0号机,B存在1号机,C存在2号机。

一切正常,直到业务增长,你需要增加第 4台 服务器。这时,计算公式变成了 hash(数据) % 4

  • 普通哈希会引发“缓存雪崩”:因为除数变了,几乎所有数据的计算结果都会改变。A、B、C很可能都不再指向原来的机器。这意味着大量的数据迁移和缓存失效,数据库压力瞬间飙升,这在生产环境是灾难性的。

  • 一致性哈希则能“优雅”应对:当增加第4台服务器时,只有大约 1/4 的数据(即 1/N) 需要重新分配。大部分数据(如A和C)依然能命中原来的机器,只有B被迁移到新机器。迁移成本被降到了最低。

底层原理:从“取模”到“环形空间”

这个神奇效果的秘密,在于它们的数据结构完全不同。

  • 普通哈希:是一个线性数组。数据的存储位置完全由取模结果决定,节点变化会直接改变整个数组的映射关系。
  • 一致性哈希:是一个哈希环(一个从0到2^32-1的闭合圆环)。服务器和数据都被哈希到这个环上,数据沿着顺时针方向找到的第一个服务器节点就是它的归宿。这种机制下,增删节点只会影响环上的一小段区域。

进阶:一致性哈希的“坑”与“填坑”

虽然一致性哈希很好,但在实际应用中,如果服务器太少(比如只有3台),它们在环上可能分布不均,导致数据“倾斜”——大部分数据都存到了某台机器上。

为了解决这个问题,业界引入了虚拟节点技术。简单说,就是给每台物理机在环上复制出几百个“分身”(虚拟节点)。这样一来,环上节点数量暴增,分布自然就均匀了,保证了负载的均衡。

场景对比:它们分别用在哪?

  • 普通哈希:适用于节点数量固定、极少变动的本地场景,比如单机缓存(如Java的HashMap)、分库分表(提前规划好库表数量)。
  • 一致性哈希:适用于节点动态变化的分布式系统,最经典的就是Redis集群分布式缓存(如Memcached),以及数据库分库分表的扩容场景。

总结一句话

  • 普通哈希像一个牢不可破的衣柜:隔板(节点)位置固定,一旦改变,所有衣服(数据)都得重新收拾。
  • 一致性哈希像一个可伸缩的圆环衣架:增减衣架上的挂钩(节点),只需挪动附近几件衣服,其他的纹丝不动。

所以,如果你的系统服务器规模需要动态扩缩容,一致性哈希是更稳妥的选择。