并发服务器:第 8 部分 - Go

Eli Bendersky 工具

摘要

本文是编写并发网络服务器系列文章的第 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">&quot;9090&quot;</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">&gt;=</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">&quot;Serving on port&quot;</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">&quot;tcp&quot;</span><span class="p">,</span><span class="w"> </span><span class="s">&quot;:&quot;</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">&quot;Error listening:&quot;</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">&quot;Error accepting connection: %v&quot;</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">&quot;connection received from&quot;</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">&quot;Error serving %v: %v&quot;</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">&quot;peer done&quot;</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

Eli Bendersky

本文是关于并发服务器系列文章的一部分,介绍了如何使用Rust实现并发网络服务器,涵盖了顺序、线程和事件驱动方法,并提供了代码示例。

C语言中的Go风格并发

Hacker News Top

一篇详细的技术文章,探讨如何在C语言中复制Go的并发模型,使用POSIX线程、互斥锁、条件变量和工作池,作为Solod转译器项目的一部分。

Goroutines 101:基础教程

Lobsters Hottest

这篇文章提供了Go语言中goroutines的基础教程,解释了它们如何简化并发以及如何有效地使用它们。

Go 中的数据竞争与内存模型

Lobsters Hottest

本文解释了数据竞争和 Go 内存模型,说明了在 goroutine 中对共享变量的非同步访问可能导致的问题,并讨论了正确的同步方法。

就用Go

Lobsters Hottest

一篇带有强烈观点的开发者文章倡导使用Go编程语言,强调其简洁的语法、强大的标准库、高效的并发模型以及单二进制部署,作为对过于复杂的现代技术栈的实用替代方案。