并发服务器:第 8 部分 - Go
摘要
本文是编写并发网络服务器系列文章的第 8 部分,重点介绍如何在 Go 中使用 goroutines 和 Go 运行时调度来实现并发。
<p>这是关于编写并发网络服务器的系列文章的第 8 部分。在本部分,我们将切换到 Go,看看它如何解决前面系列文章中描述的挑战。</p>
<p>系列中的所有文章:</p>
<ul class="simple">
<li><a class="reference external" href="https://eli.thegreenplace.net/2017/concurrent-servers-part-1-introduction/">第 1 部分 - 引言</a></li>
<li><a class="reference external" href="https://eli.thegreenplace.net/2017/concurrent-servers-part-2-threads/">第 2 部分 - 线程</a></li>
<li><a class="reference external" href="https://eli.thegreenplace.net/2017/concurrent-servers-part-3-event-driven/">第 3 部分 - 事件驱动</a></li>
<li><a class="reference external" href="https://eli.thegreenplace.net/2017/concurrent-servers-part-4-libuv/">第 4 部分 - libuv</a></li>
<li><a class="reference external" href="https://eli.thegreenplace.net/2017/concurrent-servers-part-5-redis-case-study/">第 5 部分 - Redis 案例研究</a></li>
<li><a class="reference external" href="https://eli.thegreenplace.net/2018/concurrent-servers-part-6-callbacks-promises-and-asyncawait/">第 6 部分 - 回调、承诺和异步/等待</a></li>
<li><a class="reference external" href="https://eli.thegreenplace.net/2026/concurrent-servers-part-7-rust/">第 7 部分 - Rust</a></li>
<li>第 8 部分 - Go(本部分)</li>
</ul>
<p>本文假设读者对 Go 编程语言有基本了解。</p>
<div class="section" id="sequential-state-machine-server">
<h2>顺序状态机服务器</h2>
<p>与之前一样,我们将从一个顺序服务器开始,用于 <a class="reference external" href="https://eli.thegreenplace.net/2017/concurrent-servers-part-1-introduction/">第 1 部分</a> 中介绍的基本状态机协议。</p>
<p>这是主函数:</p>
<div class="highlight"><pre><span></span><span class="kd">func</span><span class="w"> </span><span class="nx">main</span><span class="p">()</span><span class="w"> </span><span class="p">{</span><span class="w"></span>
<span class="w"> </span><span class="nx">port</span><span class="w"> </span><span class="o">:=</span><span class="w"> </span><span class="s">"9090"</span><span class="w"></span>
<span class="w"> </span><span class="k">if</span><span class="w"> </span><span class="nb">len</span><span class="p">(</span><span class="nx">os</span><span class="p">.</span><span class="nx">Args</span><span class="p">)</span><span class="w"> </span><span class="o">>=</span><span class="w"> </span><span class="mi">2</span><span class="w"> </span><span class="p">{</span><span class="w"></span>
<span class="w"> </span><span class="nx">port</span><span class="w"> </span><span class="p">=</span><span class="w"> </span><span class="nx">os</span><span class="p">.</span><span class="nx">Args</span><span class="p">[</span><span class="mi">1</span><span class="p">]</span><span class="w"></span>
<span class="w"> </span><span class="p">}</span><span class="w"></span>
<span class="w"> </span><span class="nx">log</span><span class="p">.</span><span class="nx">Println</span><span class="p">(</span><span class="s">"Serving on port"</span><span class="p">,</span><span class="w"> </span><span class="nx">port</span><span class="p">)</span><span class="w"></span>
<span class="w"> </span><span class="nx">listener</span><span class="p">,</span><span class="w"> </span><span class="nx">err</span><span class="w"> </span><span class="o">:=</span><span class="w"> </span><span class="nx">net</span><span class="p">.</span><span class="nx">Listen</span><span class="p">(</span><span class="s">"tcp"</span><span class="p">,</span><span class="w"> </span><span class="s">":"</span><span class="o">+</span><span class="nx">port</span><span class="p">)</span><span class="w"></span>
<span class="w"> </span><span class="k">if</span><span class="w"> </span><span class="nx">err</span><span class="w"> </span><span class="o">!=</span><span class="w"> </span><span class="kc">nil</span><span class="w"> </span><span class="p">{</span><span class="w"></span>
<span class="w"> </span><span class="nx">log</span><span class="p">.</span><span class="nx">Fatal</span><span class="p">(</span><span class="s">"Error listening:"</span><span class="p">,</span><span class="w"> </span><span class="nx">err</span><span class="p">)</span><span class="w"></span>
<span class="w"> </span><span class="p">}</span><span class="w"></span>
<span class="w"> </span><span class="k">defer</span><span class="w"> </span><span class="nx">listener</span><span class="p">.</span><span class="nx">Close</span><span class="p">()</span><span class="w"></span>
<span class="w"> </span><span class="k">for</span><span class="w"> </span><span class="p">{</span><span class="w"></span>
<span class="w"> </span><span class="nx">conn</span><span class="p">,</span><span class="w"> </span><span class="nx">err</span><span class="w"> </span><span class="o">:=</span><span class="w"> </span><span class="nx">listener</span><span class="p">.</span><span class="nx">Accept</span><span class="p">()</span><span class="w"></span>
<span class="w"> </span><span class="k">if</span><span class="w"> </span><span class="nx">err</span><span class="w"> </span><span class="o">!=</span><span class="w"> </span><span class="kc">nil</span><span class="w"> </span><span class="p">{</span><span class="w"></span>
<span class="w"> </span><span class="nx">log</span><span class="p">.</span><span class="nx">Printf</span><span class="p">(</span><span class="s">"Error accepting connection: %v"</span><span class="p">,</span><span class="w"> </span><span class="nx">err</span><span class="p">)</span><span class="w"></span>
<span class="w"> </span><span class="k">continue</span><span class="w"></span>
<span class="w"> </span><span class="p">}</span><span class="w"></span>
<span class="w"> </span><span class="nx">log</span><span class="p">.</span><span class="nx">Println</span><span class="p">(</span><span class="s">"connection received from"</span><span class="p">,</span><span class="w"> </span><span class="nx">conn</span><span class="p">.</span><span class="nx">RemoteAddr</span><span class="p">())</span><span class="w"></span>
<span class="w"> </span><span class="k">if</span><span class="w"> </span><span class="nx">err</span><span class="w"> </span><span class="o">:=</span><span class="w"> </span><span class="nx">server</span><span class="p">.</span><span class="nx">ServeSerialProtocol</span><span class="p">(</span><span class="nx">conn</span><span class="p">);</span><span class="w"> </span><span class="nx">err</span><span class="w"> </span><span class="o">!=</span><span class="w"> </span><span class="kc">nil</span><span class="w"> </span><span class="p">{</span><span class="w"></span>
<span class="w"> </span><span class="nx">log</span><span class="p">.</span><span class="nx">Printf</span><span class="p">(</span><span class="s">"Error serving %v: %v"</span><span class="p">,</span><span class="w"> </span><span class="nx">conn</span><span class="p">.</span><span class="nx">RemoteAddr</span><span class="p">(),</span><span class="w"> </span><span class="nx">err</span><span class="p">)</span><span class="w"></span>
<span class="w"> </span><span class="p">}</span><span class="w"> </span><span class="k">else</span><span class="w"> </span><span class="p">{</span><span class="w"></span>
<span class="w"> </span><span class="nx">log</span><span class="p">.</span><span class="nx">Println</span><span class="p">(</span><span class="s">"peer done"</span><span class="p">,</span><span class="w"> </span><span class="nx">conn</span><span class="p">.</span><span class="nx">RemoteAddr</span><span class="p">())</span><span class="w"></span>
<span class="w"> </span><span class="p">}</span><span class="w"></span>
<span class="w"> </span><span class="p">}</span><span class="w"></span>
<span class="p">}</span><span class="w"></span>
</pre></d
查看缓存全文
缓存时间: 2026/08/22 15:40
# 并发服务器:第八部分 - Go
来源:https://eli.thegreenplace.net/2026/concurrent-servers-part-8-go
这是关于编写并发网络服务器的系列文章的第八部分。在本部分中,我们将转向Go语言,看看它如何应对此前系列文章中描述的挑战。
系列文章目录:
- 第一部分 - 简介 (https://eli.thegreenplace.net/2017/concurrent-servers-part-1-introduction/)
- 第二部分 - 线程 (https://eli.thegreenplace.net/2017/concurrent-servers-part-2-threads/)
- 第三部分 - 事件驱动 (https://eli.thegreenplace.net/2017/concurrent-servers-part-3-event-driven/)
- 第四部分 - libuv (https://eli.thegreenplace.net/2017/concurrent-servers-part-4-libuv/)
- 第五部分 - Redis案例研究 (https://eli.thegreenplace.net/2017/concurrent-servers-part-5-redis-case-study/)
- 第六部分 - 回调、Promise与async/await (https://eli.thegreenplace.net/2018/concurrent-servers-part-6-callbacks-promises-and-asyncawait/)
- 第七部分 - Rust (https://eli.thegreenplace.net/2026/concurrent-servers-part-7-rust/)
- 第八部分 - Go (本文)
本文假设读者对Go编程语言有基本了解。
## 顺序状态机服务器
与之前一样,我们将从实现第一部分 (https://eli.thegreenplace.net/2017/concurrent-servers-part-1-introduction/) 中介绍的基本状态机协议的顺序服务器开始。
以下是主函数:
```
func main() {
port := "9090"
if len(os.Args) >= 2 {
port = os.Args[1]
}
log.Println("Serving on port", port)
listener, err := net.Listen("tcp", ":"+port)
if err != nil {
log.Fatal("Error listening:", err)
}
defer listener.Close()
for {
conn, err := listener.Accept()
if err != nil {
log.Printf("Error accepting connection: %v", err)
continue
}
log.Println("connection received from", conn.RemoteAddr())
if err := server.ServeSerialProtocol(conn); err != nil {
log.Printf("Error serving %v: %v", conn.RemoteAddr(), err)
} else {
log.Println("peer done", conn.RemoteAddr())
}
}
}
```
与前面部分相同,这个服务器是“无限”的;除非被显式终止,否则它会一直接受新的连接。
这是为单个客户端实现协议的函数;它接受一个`net.Conn`值,代表另一端连接着客户端的套接字:
```
type processingState int
const (
waitForMsg processingState = iota
inMsg
)
// ServeSerialProtocol 为单个TCP连接提供我们的串行协议服务。
func ServeSerialProtocol(conn net.Conn) error {
defer conn.Close()
if _, err := conn.Write([]byte{'*'}); err != nil {
return err
}
var state processingState = waitForMsg
buf := make([]byte, 1024)
for {
n, err := conn.Read(buf)
for _, b := range buf[:n] {
switch state {
case waitForMsg:
if b == '^' {
state = inMsg
}
case inMsg:
if b == '$' {
state = waitForMsg
} else {
var bb byte = byte(b) + 1
if _, err := conn.Write([]byte{bb}); err != nil {
return err
}
}
}
}
// 在处理接收到的字节及其错误之后检查错误。
if err != nil {
if errors.Is(err, io.EOF) || errors.Is(err, net.ErrClosed) {
return nil
} else {
return err
}
}
}
}
```
## 每个客户端一个goroutine
Go运行时并未直接暴露操作系统线程,而是实现了自己的M:N调度,在操作系统线程之上调度轻量级的*goroutine*。在Go中使用goroutine非常廉价——无论是在语法和开发工作量上,还是在系统资源 (https://eli.thegreenplace.net/2018/measuring-context-switching-and-memory-overheads-for-linux-threads/) 上。
这是我们的串行协议服务器的一个版本,它通过为每个客户端启动一个goroutine来并发地处理客户端。与之前示例不同的代码部分已高亮显示:
```
func main() {
port := "9090"
if len(os.Args) >= 2 {
port = os.Args[1]
}
log.Println("Serving on port", port)
listener, err := net.Listen("tcp", ":"+port)
if err != nil {
log.Fatal("Error listening:", err)
}
defer listener.Close()
for {
conn, err := listener.Accept()
if err != nil {
log.Printf("Error accepting connection: %v", err)
continue
}
log.Println("connection received from", conn.RemoteAddr())
go func() {
if err := server.ServeSerialProtocol(conn); err != nil {
log.Printf("Error serving %v: %v", conn.RemoteAddr(), err)
} else {
log.Println("peer done", conn.RemoteAddr())
}
}()
}
}
```
在这种情况下,并发修改特别简单,因为服务器是无限的;没有必要等待这些goroutine完成(因此也无需使用async.WaitGroup (https://pkg.go.dev/sync#WaitGroup))。`server.ServeSerialProtocol`的参数从外围作用域词法捕获,其返回值由外围的闭包处理。
由于goroutine非常轻量,这个服务器不太可能因为启动太多goroutine而耗尽资源;事实上,它可能会先耗尽其他资源——比如套接字的文件描述符。然而,有时限制并发程度仍然很有用——即使是在Go中,我们将在接下来的部分讨论一些实现方法。
## 使用信号量限制并发
以下是一些在Go程序中有意义限制并发程度的场景,即使goroutine启动和运行都很廉价:
- 任务可能是计算密集型的,而任何服务器的CPU容量本质上是有限的。如果太多并发的goroutine竞争有限的CPU,它们都将进展甚微。让较少的任务在合理的时间内完成可能更有意义。
- 保护可能有限的下游资源,例如并发的数据库连接或其他服务。例如,如果服务器必须为每个任务向其他服务发送请求,而这些服务又受到速率限制,那么就必须仔细管理并发性。
- 当工作由客户端指定时出于安全原因;恶意的客户端可能会使过于热衷于服务的服务过载并崩溃,使其对合法客户端不可用。
为了接下来的部分,让我们切换到第四部分 (https://eli.thegreenplace.net/2017/concurrent-servers-part-4-libuv/) 的素数测试服务器,因为它代表了一个更现实的工作负载。回顾一下:服务器接收数字,通过休眠模拟阻塞,并返回“prime”(素数)或“composite”(合数)。无限制的每个客户端一个goroutine的版本看起来几乎与之前的代码示例完全相同 (https://github.com/eliben/code-for-blog/tree/main/2026/async-socket-server-go/cmd/prime-server-unbounded),除了goroutine调用的是另一个函数:
```
go func() {
if err := server.ServePrimeProtocol(conn); err != nil {
log.Printf("Error serving %v: %v", conn.RemoteAddr(), err)
} else {
log.Println("peer done", conn.RemoteAddr())
}
}()
```
其中`ServePrimeProtocol`是[\[1\]](https://eli.thegreenplace.net/2026/concurrent-servers-part-8-go#footnote-1):
```
// ServePrimeProtocol 为单个TCP连接提供我们的素数协议服务。
func ServePrimeProtocol(conn net.Conn) error {
defer conn.Close()
buf := make([]byte, 1024)
for {
n, readerr := conn.Read(buf)
if n > 0 {
// 将读取的缓冲区解析为i64
num, err := strconv.ParseInt(strings.TrimSpace(string(buf[:n])), 10, 64)
if err != nil {
return err
}
response := "composite"
if isPrime(num, true) {
response = "prime"
}
if _, err := conn.Write([]byte(response + "\n")); err != nil {
return err
}
}
if readerr != nil {
if errors.Is(readerr, io.EOF) || errors.Is(readerr, net.ErrClosed) {
return nil
}
return readerr
}
}
}
// isPrime 如果n是素数则返回true,否则返回false。如果delay为true,
// 它会在计算前休眠n毫秒。
func isPrime(n int64, delay bool) bool {
if delay {
time.Sleep(time.Duration(n) * time.Millisecond)
}
if n < 2 {
return false
}
if n%2 == 0 {
return n == 2
}
for i := int64(3); i*i <= n; i += 2 {
if n%i == 0 {
return false
}
}
return true
}
```
在Go中限制并发最简单的方法是使用计数信号量,通过channel实现:
```
func main() {
port := "8070"
if len(os.Args) >= 2 {
port = os.Args[1]
}
log.Println("Serving on port", port)
listener, err := net.Listen("tcp", ":"+port)
if err != nil {
log.Fatal("Error listening:", err)
}
defer listener.Close()
maxConcurrency := runtime.NumCPU()
sem := make(chan struct{}, maxConcurrency)
for {
conn, err := listener.Accept()
if err != nil {
log.Printf("Error accepting connection: %v", err)
continue
}
log.Println("connection received from", conn.RemoteAddr())
// 从信号量获取一个令牌以限制并发。
sem <- struct{}{}
go func() {
// 完成服务连接后归还令牌。
defer func() { <-sem }()
if err := server.ServePrimeProtocol(conn); err != nil {
log.Printf("Error serving %v: %v", conn.RemoteAddr(), err)
} else {
log.Println("peer done", conn.RemoteAddr())
}
}()
}
}
```
channel `sem` 充当信号量;注意它是一个有界channel,具有最大容量。通过向channel发送来获取令牌,通过从channel接收来释放令牌。当channel已满时,发送操作`sem <- struct{}{}`会阻塞,直到某个其他goroutine移除一个令牌[\[2\]](https://eli.thegreenplace.net/2026/concurrent-servers-part-8-go#footnote-2)。
channel的类型是`struct{}`,表示“空的”,或“没有数据”。这在Go中是惯用法,用于那些仅用于其语义(而非发送/接收任何实际数据)的channel。
## 工作池
由于启动goroutine很廉价,并且限制并发如上所示很简单,“工作池”模式对于我们这样的服务器场景通常不是必需的。不过,它偶尔也有用处(例如,当工作线程需要跨任务维护一些非平凡的状态时),所以这里值得讨论一下。
这是使用工作池的素数测试服务器的一个变体:
```
func worker(jobs <-chan net.Conn) {
for conn := range jobs {
if err := server.ServePrimeProtocol(conn); err != nil {
log.Printf("Error serving %v: %v", conn.RemoteAddr(), err)
} else {
log.Println("peer done", conn.RemoteAddr())
}
}
}
func main() {
port := "8070"
if len(os.Args) >= 2 {
port = os.Args[1]
}
log.Println("Serving on port", port)
listener, err := net.Listen("tcp", ":"+port)
if err != nil {
log.Fatal("Error listening:", err)
}
defer listener.Close()
// 该channel是无缓冲的,因此当没有可用的工作线程接受新连接时会阻塞。
jobs := make(chan net.Conn)
maxConcurrency := runtime.NumCPU()
for i := 0; i < maxConcurrency; i++ {
go worker(jobs)
}
for {
conn, err := listener.Accept()
if err != nil {
log.Printf("Error accepting connection: %v", err)
continue
}
log.Println("connection received from", conn.RemoteAddr())
jobs <- conn
}
}
```
固定数量的worker goroutine被启动;这些goroutine都从同一个channel接收“任务”。在`Accept`循环中,每个客户端连接作为一个新任务发送到此channel,并由下一个可用的worker拾取。如前所述,通常你会看到`async.WaitGroup`来确保goroutine的干净关闭,但在我们的情况下,由于我们有一个永不退出的服务器,所以这没有必要。
## 异步?
程序员是否必须在Go中求助于异步/事件驱动编程?根据我的经验,几乎不需要。Go从下往上设计,就是为了适合大规模并发;goroutine创建成本非常低,内存占用很小,并且切换非常快速,全部发生在用户空间。我在2018年 (https://eli.thegreenplace.net/2018/measuring-context-switching-and-memory-overheads-for-linux-threads/) 进行的测量显示,切换时间约为170纳秒,而Linux上线程为1-2微秒。
此外,Go底层已经在使用像epoll这样的事件驱动循环进行I/O。等待套接字等I/O的goroutine实际上被“暂停”且不消耗资源(除了它们的小内存占用);当它们的I/O描述符就绪时,它们被Go的运行时唤醒——这与异步编程的工作方式非常相似!
也就是说,当处理*数百万*个并发流时,有些人当然会尝试通过直接在Go中进行异步编程来进一步榨取资源。我只想说这非常罕见,绝大多数用户永远不需要这样做。
## 结论
2018年,我写了一篇名为“Go精准地命中了并发的钉子” (https://eli.thegreenplace.net/2018/go-hits-the-concurrency-nail-right-on-the-head/) 的文章,经过多年的积极编码后,我完全坚持这一观点。Go对于并发程序来说极其强大且符合人体工程学;虽然其他环境竭尽全力在库中实现async/await风格的事件循环,但在Go中它已经内置于核心语言和运行时中。你想要事件驱动的I/O,拥有非常轻量级的绿色线程,同时还能执行阻塞任务而不必担心函数着色问题 (https://journal.stuffwithstuff.com/2015/02/01/what-color-is-your-function/)?Go满足了你。
## 代码
本文的所有代码可在GitHub上找到 (https://github.com/eliben/code-for-blog/tree/main/2026/async-socket-server-go)。
---
[\[1\]](https://eli.thegreenplace.net/2026/concurrent-servers-part-8-go#footnote-reference-1)仔细的读者会注意到这段代码有两个问题:(1) 该协议假设完整的数字在单次`conn.Read`调用中从套接字读取,并且没有帧结构——不同数字之间没有分隔;(2) 素数检查循环使用了`i*i`,对于大数字可能会溢出。
这些问题在C、Python、JavaScript和Rust的早期部分中的素数服务器的所有版本中都是一致的,因为我的重点是展示并发要点的最简单代码。
[\[2\]](https://eli.thegreenplace.net/2026/concurrent-servers-part-8-go#footnote-reference-2)练习:请注意,我们的信号量保护了整个`ServePrimeProtocol`,这意味着一个恶意的(或无能的)客户端可以连接并空闲而不发送任何请求,从而将我们的并发能力降低1。调整代码以将信号量移到素数计算本身周围,使只有这部分受到限制。
---
如有评论,请给我[发邮件](mailto:[email protected])。
相似文章
并发服务器:第7部分 - Rust
本文是关于并发服务器系列文章的一部分,介绍了如何使用Rust实现并发网络服务器,涵盖了顺序、线程和事件驱动方法,并提供了代码示例。
C语言中的Go风格并发
一篇详细的技术文章,探讨如何在C语言中复制Go的并发模型,使用POSIX线程、互斥锁、条件变量和工作池,作为Solod转译器项目的一部分。
Goroutines 101:基础教程
这篇文章提供了Go语言中goroutines的基础教程,解释了它们如何简化并发以及如何有效地使用它们。
Go 中的数据竞争与内存模型
本文解释了数据竞争和 Go 内存模型,说明了在 goroutine 中对共享变量的非同步访问可能导致的问题,并讨论了正确的同步方法。
就用Go
一篇带有强烈观点的开发者文章倡导使用Go编程语言,强调其简洁的语法、强大的标准库、高效的并发模型以及单二进制部署,作为对过于复杂的现代技术栈的实用替代方案。