SRS的并发模型一直挺有名气,它用协程来扛高并发,代码写起来像同步一样,但实际跑起来的调度逻辑却暗藏玄机。上个月我们的线上环境就撞到一次诡异故障:RTMP协议转FLV时,只要推流路数一多,部分HTTP拉流会卡住,最终超时。排查了很久,最后发现是多个协程回调同时操作共享数据,导致了竞态条件和死锁。修复的方法很简单,把一段“检查再更新”的流程换成原子操作,问题就消失了。这篇文章把当时的坑完整讲一遍,希望你能少走弯路。

一、问题现场:连接卡死,CPU却很闲

那天的现象是:下游播放端不断报超时,服务器端看监控,CPU占用率不到20%,内存也正常,日志没有任何报错。用 pprof 抓协程栈,看到有十几条协程全部阻塞在一个叫 bufferMu 的互斥锁上,等待锁的调用方都来自协议转换模块。也就是说,系统不是忙死,而是“等死”。

当时我们在 SRS 静态 RTMP 转 HLS/HTTP-FLV 这个场景里加了一个自定义扩展:把收到的 RTMP 音视频数据转成 FLV 后,再推给内部的播放网关。负责转换的 Worker 协程会响应网络事件回调,每个包都要调用一次 ConvertAndWrite。这个方法内部先申请一块共享缓冲区,再写入数据,最后更新缓冲区的写指针。最初代码里没有任何同步,因为开发者觉得“反正 SRS 的每个连接都有自己的协程,互不干扰”。但协议转换只有一个全局的 FLV muxer 实例,多个推流连接会同时写同一个输出流。这一下就乱了。

二、协程和回调,到底怎么跑的

在讲竞态之前,先说说 SRS 协程的基本跑法。SRS 用的协程库(State Threads)和 Go 的 goroutine 非常像,都是用户态调度。给它一个函数,它就开一个轻量级执行流,遇到 socket 读写、定时器等待会自动让出线程,等事件准备好了再恢复。因为调度是协作式的,一个协程如果不主动让出,其他协程就没有机会跑。

2.1 回调是怎么被调用的

SRS 里大量使用回调。比如你注册一个“收到 RTMP 包”的处理函数,之后每次有网络包,事件循环就会把这个函数重新拉起来执行。用 Go 模拟一下就是下面这个样子,注意这里的回调是异步的,它不一定会发生在当前协程的栈上:

技术栈:Go

// 技术栈:Go
package main

import (
	"fmt"
	"time"
)

// 收到网络包后要执行的回调
func onRtmpPacket(packet string) {
	fmt.Println("收到包:", packet)
}

// 这是SRS里一个连接协程的处理逻辑
func handleConnection() {
	// 模拟注册一个网络回调
	go func() {
		time.Sleep(1 * time.Millisecond)
		onRtmpPacket("RTMP chunk")
	}()

	// 当前协程继续干自己的事情
	time.Sleep(2 * time.Millisecond)
	fmt.Println("连接处理继续执行")
}

func main() {
	go handleConnection()
	time.Sleep(3 * time.Millisecond)
}

这里的关键点:回调函数执行的时候,主协程可能已经跑到了下一行代码,也可能正在等待其他东西。所以如果有两个连接同时触发了回调,它们对共享变量的读写顺序是不确定的。这不是线程才会有的问题,协程一样会有,因为协程之间的切换点藏在各种 IO 调用和 time.Sleep 里。

三、竞态条件和死锁的现场还原

3.1 竞态:两个人同时修改一份清单

我们的协议转换模块里有一个共享的写指针 writePos,每次转换完一个包,要把数据尾部的位置记录下来。正常逻辑是“读 writePos,写入数据,更新 writePos”。但这三步中间可能发生协程调度。例如协程 A 读到了 writePos=100,准备写入;协程 B 也读到了 writePos=100,它也准备写入。结果 A 和 B 都写到了同一个位置,后写的把先写的覆盖了,writePos 最终变成 100+len(B),而 A 的数据就丢了。

用代码还原就是这个样子,这里故意在中间让出协程,让切换更容易发生:

技术栈:Go

// 技术栈:Go
package main

import (
	"fmt"
	"runtime"
	"sync"
)

var (
	writePos int32
	output   = make([]byte, 8)
)

func convertAndWrite(data []byte) {
	// 第一步:读取当前写位置
	pos := writePos

	// 第二步:模拟做大量的协议处理,期间主动让出协程
	for i := 0; i < 10; i++ {
		runtime.Gosched()
	}

	// 第三步:把数据写到pos位置
	copy(output[pos:], data)

	// 第四步:更新写位置
	writePos = pos + int32(len(data))
	fmt.Printf("data=%q, pos=%d, newWritePos=%d\n", data, pos, writePos)
}

func main() {
	var wg sync.WaitGroup
	wg.Add(2)

	go func() {
		defer wg.Done()
		convertAndWrite([]byte("aaa"))
	}()

	go func() {
		defer wg.Done()
		convertAndWrite([]byte("bb"))
	}()

	wg.Wait()
}

运行这段代码,两个协程很可能拿到同一个 pos,然后第二次拷贝覆盖第一次的内容,控制台会打出两个相同的 pos。这种“读改写”不是原子操作,放到协同式协程里也一样会翻车。

3.2 死锁:自己锁自己

发现数据被覆盖后,我们慌忙给 convertAndWrite 加了一把互斥锁。锁是加上了,但反而引发了死锁。原因是 SRS 的协议转换流程里,回调里还会调用另一个回调,比如 ConvertAndWrite 内部会触发“缓冲进度上报”,上报函数里又要读取同一个缓冲区,于是去获取同一把锁。锁不具备可重入性,协程 A 持有锁后调用上报函数,而上报函数尝试再次持锁,瞬间就把自己锁死了。

下面这个例子完整复现了这种自死锁:

技术栈:Go

// 技术栈:Go
package main

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

var mu sync.Mutex

// 模拟SRS里两个回调互相调用
func convertAndWrite() {
	mu.Lock()
	fmt.Println("转换中...")

	// 转换完成后,调用上报函数,但这里会再次尝试获取同一把锁
	reportProgress() // 死锁!
	mu.Unlock()
}

func reportProgress() {
	mu.Lock()
	fmt.Println("上报进度")
	mu.Unlock()
}

func main() {
	go convertAndWrite()
	time.Sleep(time.Millisecond)
	fmt.Println("没有看到'上报进度',说明已经死锁")
}

运行结果只会打到“转换中...”,然后永远卡住。死锁的现场就是那么安静,CPU 几乎为零,但整个协程再也醒不过来。

四、原子操作:把“检查-修改”变成一个动作

解决竞态其实不一定要用锁。我们的写序问题本质上是“读 writePos”和“更新 writePos”之间被插入了另一段代码。如果能把这个两个动作变成一个不可分割的动作,就不会有重复序号了。Go 的 atomic 包提供了 CAS(Compare And Swap)函数:它先比较当前值是不是你期望的值,如果是就更新,如果不是就重来。整个过程由 CPU 保证原子性。

比如要让 writePos 从 cur 变成 next,可以用这个函数:

技术栈:Go

// 技术栈:Go
package main

import (
	"fmt"
	"sync/atomic"
)

var writePos int32

func tryAdvance(cur, next int32) bool {
	// 如果 writePos 还是 cur,就把它换成 next,并返回 true
	// 如果在执行过程中其他协程改了 writePos,函数返回 false
	return atomic.CompareAndSwapInt32(&writePos, cur, next)
}

func main() {
	cur := atomic.LoadInt32(&writePos)
	next := cur + 4

	if tryAdvance(cur, next) {
		fmt.Println("推进成功,新位置:", next)
	} else {
		fmt.Println("失败,说明有其他协程抢先了一步")
	}
}

CAS 是一种乐观锁思想:先假设自己拿到的值有效,更新时如果发现变了,再重新读。这样就不会出现两个协程拿到同一个序号。我们用 CAS 来修一下 convertAndWrite:

4.1 用 CAS 修复序号分配

上面 convertAndWrite 的问题在于读序号和写序号分开了。我们改成循环 CAS,保证每个协程拿到的 writePos 都是独一无二的:

技术栈:Go

// 技术栈:Go
package main

import (
	"fmt"
	"sync"
	"sync/atomic"
)

var writePos int32

// 预分配足够大的缓冲区,避免扩容竞争,这里只是为了演示
var output = make([]byte, 1024)

func convertAndWrite(data []byte) {
	for {
		// 读取当前写位置
		cur := atomic.LoadInt32(&writePos)
		next := cur + int32(len(data))

		// CAS尝试更新;失败则说明其他协程抢先,重新循环
		if atomic.CompareAndSwapInt32(&writePos, cur, next) {
			// 这个位置是当前协程独享的,可以安全写入
			copy(output[cur:next], data)
			fmt.Printf("data=%q 写到了位置[%d,%d)\n", data, cur, next)
			return
		}
	}
}

func main() {
	var wg sync.WaitGroup
	wg.Add(2)

	go func() {
		defer wg.Done()
		convertAndWrite([]byte("aaa"))
	}()

	go func() {
		defer wg.Done()
		convertAndWrite([]byte("bb"))
	}()

	wg.Wait()
}

注意,这里 CAS 只保护 writePos 的分配,output 本身仍然是共享内存,但不同协程拿到的位置不重叠,所以 copy 不会碰到对方的字节。如果以后要支持“追加数据”或“重写前一个包”,那就不能只依赖 CAS 了,需要配合更完整的同步方案。

4.2 原子操作的好处和边界

原子操作没有锁,所以不会造成持锁等锁的死锁,性能也高。但它的边界很清楚:只能保护一个内存地址上的单个操作。你要是想保护“先检查缓冲区剩余空间,再扩容,再写入”这样一整段逻辑,CAS 一个人是忙不过来的。另外 CAS 循环在高竞争下会一直重试,占用 CPU,所以不适合用在临界区特别长的场景。

五、彻底避免死锁的配套调整

修复了序号以后,我们把之前的互斥锁也去掉了。但 SRS 的回调链仍然存在,还是要小心“持锁调用回调”这种行为。我们重新设计了转换流程,所有回调只做两件事:把数据放进队列,以及更新一个原子计数器。真正写缓冲区的工作,统一交给一个专门的协程去做。这样就不会出现一个协程既拿着锁又要去请求另一个协程手里的锁。

代码示例如下,这里用 channel 充当队列:

技术栈:Go

// 技术栈:Go
package main

import (
	"fmt"
	"sync"
	"sync/atomic"
)

type flvPacket struct {
	data       []byte
	sequenceId uint64
}

var packetCh = make(chan flvPacket, 1024)
var seq uint64

func onRtmpConverted(data []byte) {
	// 原子分配序列号,不会重复
	id := atomic.AddUint64(&seq, 1)
	// 只把数据丢进队列,立即返回
	packetCh <- flvPacket{data: append([]byte(nil), data...), sequenceId: id}
}

// 唯一的FLV缓冲写协程
func flvWriter() {
	for pkt := range packetCh {
		// 这里可以放心地写缓冲区,因为只有一个协程在写
		fmt.Printf("写FLV包 seq=%d data=%q\n", pkt.sequenceId, pkt.data)
	}
}

func main() {
	go flvWriter()

	var wg sync.WaitGroup
	for i := 0; i < 10; i++ {
		wg.Add(1)
		go func(i int) {
			defer wg.Done()
			onRtmpConverted([]byte{byte('A' + i)})
		}(i)
	}
	wg.Wait()
	close(packetCh)
}

这个方案让“分配序列号”保持原子,让“写缓冲区”变得串行,从根本上消除了竞态。同时由于没有持锁,也就没有死锁的土壤。注意 channel 本身也有锁,但它内部实现不会和业务回调互相嵌套,用起来安全很多。

六、应用场景、优缺点和注意事项

6.1 这个方案适合哪些场景

如果你在多个协程回调里需要生成唯一递增序号,比如给 HTTP-FLV 包编号、给 RTMP 消息分配时间戳等,用原子操作最合适。当多个协程回调要共享比较大的数据结构(比如缓冲区、连接池)时,不要用原子操作硬扛,建议像上面的 flvWriter 一样,把共享资源收敛到一个 goroutine 里,用 channel 串行访问。

6.2 技术优缺点

原子操作优点:开销小,不需要上下文切换,不会死锁。缺点:只能守护单个变量;无法解决多步复合逻辑的竞态;高竞争下 CAS 会空转发热。互斥锁优点:能保护任意长的临界区;缺点:一不小心就死锁,尤其不可重入,且持锁期间让出协程会导致其他协程长时间等待。channel 优点:代码逻辑清晰,天然串行化;缺点:需要额外管理队列,当队列满时会阻塞发送方,实质上是把锁转移到了 channel 内部。

6.3 注意事项

第一,原子操作不能替代锁。像“读取缓冲区长度,检查够不够,不够就扩容,然后再写入”这种多步组合,必须使用锁或 channel。第二,SRS 的协程是协作式调度,如果一个协程在持锁期间调用 sleep 或 IO,相当于把锁霸占了,其他协程只能干等。所以任何时候都不要在持锁状态下调用回调函数。第三,写并发代码要开数据竞争检测。Go 里有专门的 race 检测器,运行或测试时加上 -race 参数,就能自动抓出数据竞争。如果跑出“DATA RACE”的报错,一定要警惕,别用“概率很小”来安慰自己。

七、总结

回到最初的故障,问题并不是 SRS 本身有 bug,而是我们在协议转换模块里忽略了协程回调的并发性。通过把写序号的“读改写”换成 atomic.CAS,再把真正写缓冲区的工作放进单独协程,竞态条件消失了,死锁也自然没有了。这个案例告诉我们:在多回调环境下,同步方案一定要想清楚临界区是谁、锁会不会被递归持有、持锁时会不会让出协程。原子操作是很好的轻量武器,但不是万能药。理解协程调度,比记住多少个 API 更重要。