做资产扫描器,真正棘手的地方往往不是探测逻辑本身,而是并发控制。资产扫描的本质,是在授权范围内对一批 IP 做端口开放情况盘点,你面对的组合通常是“几万个 IP × 几百个端口”,乘出来就是上千万级别的探测任务。如果探测脚本一写完就急着全量压出去,先崩掉的基本不会是目标,而是你自己的网络栈、文件描述符,以及目标侧防火墙的告警规则。去年我重构公司内部的资产扫描引擎时,最大的体会就是:Go 的 goroutine 确实很轻,但只有把协程池和限流机制设计好了,扫描器才真正“扫得完、扫得准、扫得稳”。这篇文章就围绕我在这个项目里的协程池与限流实现展开,讲清楚为什么这么设计,以及线上跑挂之后踩过的那些坑。
1. 为什么资产扫描器离不开高并发设计
1.1 扫描任务量级远超直觉
先算一笔账。假设我们要对一个 C 段做资产发现,就是常见的 254 个可用 IP,再挑 20 个常见端口去探测,组合数是 254 × 20 = 5080 次 TCP 连接。这个量级单线程跑也就几分钟,不痛不痒。但真实场景往往是上级给一个 /16 的授权范围,65 个 C 段,接近 65534 个 IP;常见端口从 21、22、80、443 往上数,挑 100 个很保守。组合数直接变成 65534 × 100 = 655 万次探测。
这不是我夸张,做外网资产梳理时这属于起步量级。很多内部系统还涉及多个 CIDR 白名单、旁路资产、历史遗留地址段,几千万次连接很常见。单线程做这 655 万次探测是什么概念?一次 TCP 握手在公网平均 RTT 取 200ms,串行执行就是 0.2 秒一次,655 万次约等于 131 万秒,也就是差不多 15 天。授权测试的时间窗口通常只有几个晚上,这显然不可接受。只有高并发才能把时间压缩到小时级。
1.2 扫描本质是网络 IO 密集,CPU 几乎全程在等
端口探测的过程,本质是一个“发请求、等响应”的模型。客户端发出 SYN,然后等目标返回 SYN-ACK;目标没监听这个端口,会回 RST;目标有防火墙策略或网络丢包,就什么都不回,客户端只能一直等到超时。整个过程 CPU 几乎不参与计算,纯粹是内核协议栈在等待网络事件。
所以扫描器是一个典型的 IO 密集型负载。每个 goroutine 大部分时间都阻塞在 DialTimeout 上,CPU 占用率平时可能只有 10% 出头。这种情况下,并发才是吞吐量的核心。你可以把每个探测任务想成餐厅里的一桌客人:服务员下单后不需要在后厨窗口傻等着那道菜出锅,而是可以先去服务下一桌,过一会儿再把菜端回来。多协程做的事情其实就是这个——同一个服务员同时“挂起”多张订单,等待期间继续接新单。
1.3 为什么选择 Go 而不是其他语言
扫描器的高并发实现,我之前用 Python 写过一版,也参考过 Java 的实现方式,最后整体迁移到 Go,核心原因是 Go 的并发模型和这个场景太匹配了。
Python 的 asyncio 能解决 IO 密集问题,但你要写大量 await,链路一旦拉长,回调式的心智负担会非常重;加上 GIL 的存在,多线程在网络 IO 场景虽然有优势,但线程切换和内存开销并不理想。Java 可以做,Thread 本身足够成熟,NIO 生态也很强,但线程默认栈空间是 MB 级别,起十万个线程不现实,要上 Netty 这类框架又带来额外的学习成本。
Go 的 goroutine 初始栈只有 2KB,可以动态增长,配合 GMP 调度器,单机起几万甚至几十万个 goroutine 都没压力。更重要的是 channel + select 这套通信原语,让“任务分发、并发执行、结果回收”这些扫描器天然需要的工作流,代码写起来非常直白,不需要额外的线程池框架。写完调度逻辑,剩下的精力可以全部放在限流策略和稳定性上。
2. 协程池设计:管住并发边界比一味加速更重要
2.1 为什么不能无脑给每个任务开一个 goroutine
goroutine 很轻,不代表可以无限制。最常见的翻车现场,就是直接 for 循环里 go func,一次性把 655 万个探测任务全部塞进调度器。goroutine 数量确实能撑住,但每个探测任务都要创建一个 socket、占用一个本地临时端口、消耗一个文件描述符。这个压力不会平均分配到时间轴上,而是瞬间爆发。
文件描述符的默认 ulimit 通常是 1024,一秒钟创建几千个 socket 就触顶了。即使把 ulimit 调到 65535,本地源端口范围 net.ipv4.ip_local_port_range 默认 32768 到 60999,总共也就 28232 个可用端口。每个探测连接结束后会进入 TIME_WAIT 状态,在 2MSL 时间内端口不可复用,大概 60 秒。按 500 QPS 算,60 秒内会产生 30000 个 TIME_WAIT,已经逼近源端口上限。所以并发不是越高越好,必须用协程池把并发数限制在可控范围内。
2.2 协程池的三大核心要素
一个可用的协程池需要三个东西:任务队列、固定数量的 worker、结果回收通道。
任务队列我用带缓冲的 channel 实现,生产端往里提交 ScanTask,worker 从中消费。固定 worker 数量是限流的第一道闸门,核心逻辑就一个 for + select 循环,监听上下文取消和任务通道两个事件。
type Pool struct { tasks chan ScanTask results chan ScanResult wg sync.WaitGroup } func (p *Pool) Start(ctx context.Context) { for i := 0; i < p.workerCount; i++ { p.wg.Add(1) go func() { defer p.wg.Done() for { select { case <-ctx.Done(): return case task, ok := <-p.tasks: if !ok { return } p.runTask(ctx, task) } } }() } }这里的 select 是关键。只监听 tasks 而不监听 ctx 的做法,在扫描器里很容易出问题:程序收到退出信号后,worker 还会继续消费队列里的残留任务,如果队列高达几万条,扫完全部任务可能还要很久。加上 ctx.Done() 分支后,取消信号一到,worker 立刻停止拉取新任务,退出循环,整个扫描器就能快速收敛。
2.3 worker 数量怎么定
worker 数量最忌讳的是按 CPU 核数乘一个倍数来拍板。扫描器的瓶颈在网络往返和远程目标处理能力,不在 CPU。与其套公式,不如从三个物理约束反推:
第一是本地源端口和文件描述符上限。第二是出口带宽,TCP 探测每个连接至少一个包出去、一个包回来,算一下带宽能扛多少 QPS。第三是目标侧能承受的速率,这个往往最致命,目标抗不住你就会看到大批超时。
我用过的参考区间大致是这样:
- 内网资产盘点,RTT 低、带宽充裕:worker 数可以开 500 到 1000,对应 QPS 可以到 3000 以上。
- 公网授权资产发现:worker 数建议 200 到 500,QPS 控制在 300 到 800。
- 大规模互联网侧探测,追求稳定和低调:worker 100 到 200,QPS 100 到 300。
以上只是起始参考,最终值要靠实际扫描时的丢包率和误报率去调,没有一组参数能适配所有网络环境。
2.4 任务队列缓冲大小同样重要
任务队列缓冲大小,很多人直接填一个大数,觉得越大提交越快。其实缓冲过大会有两个问题:一是占用内存,ScanTask 本身不大,但几十万任务堆在 channel 里,GC 压力会上升;二是任务积压会造成数据陈旧,等你扫到一个端口的时候,扫描动作和扫描结果之间已经隔了很长时间,对实时资产发现不友好。
一般我会按 workerCount × 100 左右的量级来设置临时缓冲。比如 200 个 worker,队列给 20000。这个量级既能承受生产端一次性把一批任务灌进来,又不会积压到数据失真。更精细的做法是,生产端生成完一个批次再提交一个批次,而不是一次性把千万任务全塞进队列。
3. 限流机制设计:控制节奏才能扫得更稳
3.1 限流不只是为了保护目标,更是为了扫描结果准确
很多人第一次写扫描器会觉得限流是“慢一点”,是无奈之举。但实际跑过大规模扫描之后你就知道,限流的直接价值是保证结果可信。
当并发太猛的时候,目标机器或者中间的防火墙设备会进入防御状态:丢包、限速、甚至直接 RST。丢包导致的现象是,原来可能开放的端口,因为 SYN-ACK 回不来,客户端超时,被误判为“关闭”。一场扫描跑完,如果不去压一下速率,你看到的结果里混着大量假阴性,资产梳理的准确率根本没法保证。
另外,限流也是保护扫描器自身。前面算过,500 QPS 持续 60 秒就能攒下 3 万个 TIME_WAIT 连接,如果不限流而疯狂发包,本地网络栈先吃不消。你可以把限流理解成给水管装一个减压阀:水压太大时,先爆的往往是自家管道,而不是远端水厂。
3.2 令牌桶、漏桶与 Go 的 rate.Limiter
限流算法里最常用的就是令牌桶和漏桶。漏桶确保流量以固定速率输出,适合必须严格平滑的场景;令牌桶则允许一定程度的突发,桶里有攒下的令牌时,可以短时间打出一波流量,但总体速率仍然受控。
资产扫描器最适合的是令牌桶。因为扫描任务天然有批次性,比如刚加载完一个 IP 段,瞬间需要发出大量探测请求,如果被漏桶压成匀速,整批任务会被拖慢很多;令牌桶的突发额度能吸收这类锯齿。Go 官方库 golang.org/x/time/rate 提供了现成的令牌桶实现,使用很简单:
limiter := rate.NewLimiter(rate.Limit(500), 100) err := limiter.Wait(ctx)第一个参数是每秒补充速率,也就是稳态 QPS 上限;第二个参数是桶容量,允许瞬时突发的数量。实际扫描时我会把 Wait 放在 worker 消费任务之后、发起 TCP 探测之前。这样限流等待不会拖住生产端,worker 会自然阻塞在 Wait 上,等令牌到位再开始下一个探测。
3.3 三层限流架构
单做一层全局 QPS 限流已经能解决大部分问题,但线上跑了几天你就会发现不够。至少需要三层。
第一层是全局限流,就是上面说的 rate.Limiter,控制整台扫描器对外发出探测的总速率。
第二层是单目标 IP 限流。全局 QPS 限到 500,如果这 500 个请求全部集中在同一个目标 IP 上,对那个目标来说就是一瞬间 500 个并发连接。很多设备直接就把你封了。解决办法是给每个目标 IP 建一个独立的令牌桶:
type perHostLimiter struct { mu sync.Mutex lrate rate.Limit burst int lmap map[string]*rate.Limiter } func (h *perHostLimiter) get(ip string) *rate.Limiter { h.mu.Lock() defer h.mu.Unlock() if l, ok := h.lmap[ip]; ok { return l } l := rate.NewLimiter(h.lrate, h.burst) h.lmap[ip] = l return l }每个任务执行前先调 get(task.IP).Wait(ctx),保证对同一个目标 IP 的探测速率是受限的。需要注意 map 会越攒越大,扫描几十万 IP 后要定期清理过期条目,或者按时间窗口重建 map。
第三层是动态背压。再好的静态配置,也扛不住网络环境突变。比如扫描过程中目标侧突然出现丢包,探测耗时从 200ms 漂到 2 秒,任务队列开始积压,这时候如果还按原来的 QPS 发包,只会让情况更糟。我用一个监控 goroutine 周期采样任务队列长度,超过阈值就把全局限流器的速率降下来,队列恢复再升回去:
func (p *Pool) adaptiveMonitor(ctx context.Context, threshold int) { ticker := time.NewTicker(200 * time.Millisecond) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: qlen := len(p.tasks) switch { case qlen > threshold: p.limiter.SetLimit(p.limiter.Limit() / 2) case qlen < threshold/3: p.limiter.SetLimit(rate.Limit(p.maxQPS)) } } } }这里要留滞回区间,别在阈值附近反复升降速,否则限流器会抖得很厉害。阈值一般取队列容量的 70%,恢复点设在 30% 以下,两个点之间保持当前速率不动。
4. 完整实现:把协程池与限流组装到扫描器里
4.1 模块划分
我把扫描器调度部分拆成了四个组件:任务定义、协程池、限流器、结果处理。协程池不关心探测逻辑,只负责调度;探测函数通过注入的方式传给 worker。限流器作为 Pool 的内部依赖,对外只暴露 Submit 和 Results 两个方法。
整个流程是这样的:生产端生成 ScanTask,通过 Submit 提交到 tasks channel;固定数量的 worker 消费任务,每个任务执行前先向全局限流器请求令牌,再执行探测,探测结果写入 results channel;另一个 goroutine 持续消费 results,把开放端口输出或写入存储。任务全部提交完后关闭 tasks,worker 全部退出后关闭 results,主循环自然结束。
4.2 核心代码实现
下面是一份简化但可以直接跑通的核心代码,去掉了具体的探测指纹,只保留 TCP 端口状态判断逻辑:
package main import ( "context" "fmt" "net" "os" "os/signal" "strconv" "sync" "syscall" "time" "golang.org/x/time/rate" ) type ScanTask struct { IP string Port int } type ScanResult struct { Task ScanTask Open bool Duration time.Duration Err error } type PoolConfig struct { WorkerCount int QueueSize int QPS int Burst int DialTimeout time.Duration } type Pool struct { cfg PoolConfig tasks chan ScanTask results chan ScanResult limiter *rate.Limiter wg sync.WaitGroup } func NewPool(cfg PoolConfig) *Pool { if cfg.DialTimeout <= 0 { cfg.DialTimeout = 2 * time.Second } return &Pool{ cfg: cfg, tasks: make(chan ScanTask, cfg.QueueSize), results: make(chan ScanResult, cfg.QueueSize), limiter: rate.NewLimiter(rate.Limit(cfg.QPS), cfg.Burst), } } func (p *Pool) Start(ctx context.Context) { for i := 0; i < p.cfg.WorkerCount; i++ { p.wg.Add(1) go p.worker(ctx) } go func() { p.wg.Wait() close(p.results) }() } func (p *Pool) worker(ctx context.Context) { defer p.wg.Done() for { select { case <-ctx.Done(): return case task, ok := <-p.tasks: if !ok { return } if err := p.limiter.Wait(ctx); err != nil { return } p.sendResult(ctx, p.scan(task)) } } } func (p *Pool) scan(task ScanTask) ScanResult { start := time.Now() conn, err := net.DialTimeout("tcp", net.JoinHostPort(task.IP, strconv.Itoa(task.Port)), p.cfg.DialTimeout, ) if err != nil { return ScanResult{Task: task, Err: err, Duration: time.Since(start)} } conn.Close() return ScanResult{Task: task, Open: true, Duration: time.Since(start)} } func (p *Pool) sendResult(ctx context.Context, res ScanResult) { select { case p.results <- res: case <-ctx.Done(): } } func (p *Pool) Submit(ctx context.Context, task ScanTask) error { select { case p.tasks <- task: return nil case <-ctx.Done(): return ctx.Err() } } func (p *Pool) Results() <-chan ScanResult { return p.results } func main() { ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) defer stop() pool := NewPool(PoolConfig{ WorkerCount: 200, QueueSize: 20000, QPS: 500, Burst: 100, DialTimeout: 2 * time.Second, }) pool.Start(ctx) go func() { defer close(pool.tasks) ips := []string{"198.51.100.1", "198.51.100.2", "198.51.100.3"} ports := []int{21, 22, 80, 443, 3389, 8080} for _, ip := range ips { for _, port := range ports { select { case <-ctx.Done(): return default: } if err := pool.Submit(ctx, ScanTask{IP: ip, Port: port}); err != nil { return } } } }() openCount := 0 for res := range pool.Results() { if res.Open { openCount++ fmt.Printf("open: %s:%d (%s)\n", res.Task.IP, res.Task.Port, res.Duration.Round(time.Millisecond)) } } fmt.Println("done, open ports:", openCount) }这里有几个细节值得说明。sendResult 用 select 包裹结果发送,是因为极端情况下消费者已经因为外部原因不再接收结果,worker 不能无限阻塞在发送上。limiter.Wait 放在任务取出之后,而不是任务取出之前,是为了让限流等待只影响单个 worker,不阻塞整个队列。
4.3 动态背压与单目标限流接入
上面这份代码只有全局限流,实际部署时我会加一节中的 adaptiveMonitor 和 perHostLimiter。接入点分别在两个位置:adaptiveMonitor 在 Pool.Start 里启动一个独立 goroutine 即可;单目标限流则是在 scan 函数执行前加一道 Wait,建议放在全局限流的 Wait 之后。
可以这样改 scan 之前的逻辑:
if err := hostLimiter.get(task.IP).Wait(ctx); err != nil { return }注意单目标限流的速率要跟全局限流协调。比如全局 QPS 是 500,目标是一个 /24 的 254 个 IP,那每个 IP 平均带宽只有每秒 2 次左右。所以单目标速率我会独立设置,通常按单个目标 IP 每秒不超过 20 个新连接来限制,具体看目标服务的抗性。
4.4 参数配置建议
最终我把一组常用参数整理成了表格,方便不同场景直接起跑:
| 场景 | WorkerCount | QueueSize | QPS | Burst | DialTimeout |
|---|---|---|---|---|---|
| 内网资产盘点 | 500 | 50000 | 3000 | 500 | 500ms |
| 公网授权资产发现 | 200 | 20000 | 500 | 100 | 2s |
| 大范围低激进扫描 | 100 | 10000 | 150 | 50 | 3s |
参数不是固定的,跑第一轮时如果看到失败率超过 5%,优先降 QPS;如果队列一直打不满导致吞吐低,可以加 workerCount 或加大 Burst。扫描器调参的本质是找当前网络环境下的平衡点。
5. 线上实战:踩过的坑与排查方法
5.1 端口误报与重试机制
第一个坑是误报。第一版扫描器我把连接超时视为“端口关闭”,结果线上跑了一轮,一个明明对外开放的 Web 服务被判定成关闭了。排查后发现不是代码逻辑错了,而是对应网络路径上的防火墙设备在高并发下开始丢包,导致 SYN-ACK 回不来。
处理方式分两层。第一层是结果分级,不能只分“开放/关闭”,要加入“未知/被过滤”状态。TCP 连接得到 connection refused,基本可以确定端口关闭;得到 timeout,则不一定是关闭,可能是被丢包,也可能是目标根本没监听但防火墙静默丢弃。第二层是重试策略,只对超时类错误做重试,最多两轮,重试时把探测间隔拉长。被判定为关闭的 connection refused 不需要重试。
5.2 本地端口耗尽与 TIME_WAIT
扫描跑到 10 分钟左右,突然大批任务报错,错误信息是 cannot assign requested address。这是典型的本地源端口耗尽。排查第一步就是看 netstat 里的 TIME_WAIT 数量,我当时观察到 2 万个以上,临时端口范围已经被占满。
解决三板斧:一是把扫描 QPS 降下来,让 TIME_WAIT 生成速率小于回收速率;二是调整内核参数 net.ipv4.tcp_tw_reuse=1,并确认 net.ipv4.tcp_timestamps=1,让内核在安全前提下复用 TIME_WAIT 状态的连接;三是把 net.ipv4.ip_local_port_range 从默认的 32768 60999 扩大到 1024 65535,给源端口多腾出几千个名额。注意 tcp_tw_reuse 只对出站连接生效,恰好适合扫描器这种主动发起探测的场景。
5.3 goroutine 泄漏的典型场景
协程池写完后,我用 runtime.NumGoroutine 打印 goroutine 数量,发现扫描结束后数量居然没有回落到基线。排查过程比较折腾,最后定位到是一个 worker 里的 Wait(ctx) 传错了 context。
我一开始在 worker 的 select 里监听的是全局 ctx,但 limiter.Wait 用的是另一个 WithCancel 产生的子 ctx,子 ctx 在某种条件下没有被取消,导致部分 worker 一直阻塞在 Wait 上。后来统一所有 goroutine 都用同一个根 ctx,并且在 sendResult 里也加了 ctx 检查,goroutine 数量才恢复正常。排查 goroutine 泄漏的标准套路是:先打印 NumGoroutine 和 pprof.Lookup("goroutine"),再用 go tool pprof 抓 goroutine 堆栈,看阻塞点到底停在哪个函数。
5.4 限流等待会吃掉整体超时
还有一个容易忽略的问题,限流等待本身会消耗时间。当全局 QPS 调得比较低,而任务量很大的时候,一个任务可能要在 limiter.Wait 上等很久。如果整轮扫描设了 30 分钟的 WithTimeout,这 30 分钟里有一部分时间其实是在等令牌,真正用于探测的时间会缩水。
我的处理方式是给整个扫描任务单独一个 deadline,而不是让每个探测的超时来承担整体时间约束。ctx.WithTimeout 控制整体扫描周期,DialTimeout 控制单次 TCP 握手时间,两者分开设置。同时限流等待的时间也要计入整体预算,算 QPS 的时候预留余量,不要让任务提交速度长期逼近限流上限。
5.5 常见问题速查表
整理一份实际问题排查表,给后来的人直接对照:
| 症状 | 可能原因 | 排查与解决 |
|---|---|---|
| 大量 cannot assign requested address | 本地源端口耗尽 | 调整 ip_local_port_range,开启 tcp_tw_reuse,降低 QPS |
| 端口大量误判为关闭 | 目标侧防火墙丢包、扫描速率过高 | 降低全局 QPS,超时类错误单独重试,加入结果分级 |
| 扫描结束后 goroutine 数不回降 | context 未正确传递或未取消 | 打印 pprof goroutine 堆栈,检查阻塞点;统一根 ctx |
| 队列持续积压、吞吐上不去 | worker 数不足或单目标限流过严 | 增加 worker,检查单目标限流速率,观察 RTT 和失败率 |
| 单目标被防火墙断开 | 对同一目标并发过高 | 启用 perHostLimiter,限制单 IP 每秒连接数 |
| 扫描过程自身 CPU 不高但很卡 | 文件描述符或内存受限 | ulimit 调整,检查 channel 缓冲是否过大 |
这套架构在我线上跑过几个月,最大的感受是协程池和限流不是“性能优化”,而是扫描器的生存底线。没有它们,哪怕探测逻辑写得再漂亮,跑一轮大规模任务也会被各种资源瓶颈和误报淹死。先把任务量级算清楚,再定并发边界,最后用限流把速率收住,这三步做到位,扫描器的稳定性和准确率自然就上来了。最后提醒一句:所有扫描动作都要在授权范围内进行,资产发现也好、端口判断也好,跑之前确认目标范围合规,这是底线。