Golang 网络轮询器 netpoll:把阻塞 IO 变成事件通知
网络轮询器 netpoll:把阻塞 IO 变成事件通知
Go 的 net/http 能轻松支撑数十万并发连接,秘密之一就是 netpoll。它让 goroutine 在等待网络 IO 时,不会阻塞底层的 OS 线程,而是把 socket 注册到操作系统提供的多路复用机制上,事件就绪后再唤醒对应的 G。
一、为什么需要 netpoll
假设一个 goroutine 执行 conn.Read(),如果 OS 线程真的阻塞在 read 系统调用上,那其他 goroutine 就没法用这个 M 了。Go 的做法是:
- G 需要等待网络事件时,调用
gopark挂起自己。 - socket 被注册到 epoll/kqueue。
- 某个 M 在调度循环中调用
netpoll(),获取就绪的 socket 列表。 - 把对应的 G 标记为可运行,放回 LRQ。
这样,一个 M 就能同时照看成千上万个等待中的网络连接。
二、跨平台实现
Go runtime 在不同操作系统上使用不同的多路复用机制:
| 平台 | 机制 |
|---|---|
| Linux | epoll |
| macOS / BSD | kqueue |
| Windows | IOCP |
接口被封装成统一的形式:netpollinit、netpollopen、netpollclose、netpoll。
三、pollDesc 与文件描述符的绑定
每个网络 socket 都会关联一个 pollDesc 结构,核心字段包括:
type pollDesc struct {
link *pollDesc // 链表指针
fd uintptr // 文件描述符
rg uintptr // 等待读事件的 G 或 gList
wg uintptr // 等待写事件的 G 或 gList
user uint32 // 用户事件掩码
}
当 G 调用 netpollblock(pd, mode) 时,会把当前 G 挂到 rg(读)或 wg(写)上,然后 gopark 让出 CPU。事件就绪后,netpollready 唤醒对应的 G。
四、netpoll 什么时候被调用
调度循环会在合适的时机调用 netpoll():
schedule()找不到可运行的 G 时。sysmon系统监控线程检测到 netpoll 结果需要处理时。- 显式的网络相关阻塞/唤醒路径中。
netpoll() 有一个 timeout 参数。如果为 0,表示非阻塞轮询一次;如果大于 0,可以让出 CPU 等待事件。
五、从用户代码到 netpoll 的路径
以 net.TCPConn.Read 为例,简化路径是:
conn.Read(buf)
└── poll.FD.Read
└── fd.pd.waitRead()
└── netpollblock(&fd.pd, 'r', true)
└── gopark(...) // G 被挂起,M 继续调度其他 G
当 socket 有数据可读时,epoll/kqueue 通知 runtime,netpoll 遍历就绪事件,调用 netpollready 把 G 状态改回 _Grunnable。
六、代码实践:模拟网络 IO 事件等待
下面用 net.Pipe() 创建一对内存 pipe,演示 goroutine 等待另一端写入。虽然 pipe 不完全等于 socket,但它同样会触发 netpoll 路径。
package main
import (
"fmt"
"io"
"net"
"sync"
"time"
)
func main() {
c1, c2 := net.Pipe()
var wg sync.WaitGroup
wg.Add(2)
go func() {
defer wg.Done()
buf := make([]byte, 128)
fmt.Println("reader: waiting for data at", time.Now().Format("15:04:05.000"))
n, err := c1.Read(buf)
if err != nil && err != io.EOF {
fmt.Println("reader error:", err)
return
}
fmt.Printf("reader: got %q at %s\n", string(buf[:n]), time.Now().Format("15:04:05.000"))
}()
go func() {
defer wg.Done()
time.Sleep(300 * time.Millisecond)
fmt.Println("writer: sending at", time.Now().Format("15:04:05.000"))
_, err := c2.Write([]byte("hello netpoll"))
if err != nil {
fmt.Println("writer error:", err)
}
c2.Close()
}()
wg.Wait()
}
输出显示 reader 在 Read 处等待 300ms,直到 writer 写入数据。等待期间,reader 所在的 G 被 park,M 可以去执行其他 goroutine。
七、用 TCP Listener 观察 Accept 等待
package main
import (
"fmt"
"net"
"time"
)
func main() {
ln, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
panic(err)
}
defer ln.Close()
addr := ln.Addr().String()
fmt.Println("listen on", addr)
// 延迟 200ms 再连接
go func() {
time.Sleep(200 * time.Millisecond)
conn, err := net.Dial("tcp", addr)
if err != nil {
fmt.Println("dial error:", err)
return
}
defer conn.Close()
fmt.Println("client: connected at", time.Now().Format("15:04:05.000"))
}()
fmt.Println("server: waiting for accept at", time.Now().Format("15:04:05.000"))
conn, err := ln.Accept()
if err != nil {
panic(err)
}
defer conn.Close()
fmt.Println("server: accepted at", time.Now().Format("15:04:05.000"))
}
Accept 也是一个会走 netpoll 的阻塞调用。在等待连接期间,G 被 park,不会占用 M。
八、大量连接下的 M 数量
正是因为 netpoll,Go 程序即使有几十万个等待 IO 的 goroutine,M 的数量也可以保持在一个合理范围。没有被 IO 阻塞的 M 可以去执行其他 G,而 epoll/kqueue 负责在事件就绪时通知 runtime。
九、逐条解读
net.Pipe()创建一对全双工连接,底层是匿名 socket。c1.Read()阻塞时,当前 G 调用gopark挂起,socket 被注册到 netpoll。- 300ms 后
c2.Write()写入数据,触发 epoll/kqueue 就绪事件。 netpoll在调度循环中被调用,拿到就绪的 G,放回 LRQ。ln.Accept()同样走 netpoll 路径,等待连接事件。- 这种设计让少量 M 能支撑大量并发网络连接。
十、小结
| 概念 | 一句话总结 |
|---|---|
| netpoll | Go runtime 的网络事件轮询器,把阻塞 IO 变成异步事件 |
| epoll/kqueue/IOCP | 操作系统提供的多路复用机制 |
| pollDesc | socket 与等待 G 之间的关联结构 |
| gopark/goready | 挂起和唤醒 goroutine 的核心原语 |
| netpollblock | 把 G 注册到 pollDesc 并 park |
| netpollready | 事件就绪后唤醒 G |
openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构
更多推荐
所有评论(0)