
本文详解一个基于 go 协程与通道(channel)构建的请求调度系统,通过工作池(worker pool)模型将 http 请求排队、分发与异步执行,避免并发失控,适用于高并发写入场景(如日志收集、文件存储等)。
本文详解一个基于 go 协程与通道(channel)构建的请求调度系统,通过工作池(worker pool)模型将 http 请求排队、分发与异步执行,避免并发失控,适用于高并发写入场景(如日志收集、文件存储等)。
该代码实现了一个典型的生产者-消费者模式,核心目标是:将大量 HTTP POST 请求(模拟 AJAX 请求)安全、可控地排队,并由固定数量的工作协程(Worker)依次消费、延迟处理并持久化。整个流程不依赖第三方库,完全基于 Go 原生并发原语(goroutine + channel)构建。
✅ 核心组件解析
- RequestQueue:带缓冲的通道(make(chan Request, 1024)),作为请求生产者(HTTP handler)与调度器之间的缓冲区。写入不会阻塞(除非满),实现了请求“节流”和瞬时峰值保护。
- WorkerQueue:带缓冲的通道(chan chan Request),用于分发任务。它不直接传递请求,而是传递每个 Worker 的专属输入通道(worker.Request),即“把空篮子交给调度员”。
-
Worker 结构体:每个 Worker 持有一个专属的 Request chan Request,并启动一个长期运行的 goroutine,在其中循环:
- 将自身 Request 通道“注册”到 WorkerQueue(w.WorkerQueue
- 阻塞等待该通道接收请求(case request :=
- 执行业务逻辑(如 time.Sleep(request.Delay) 和 writeToFile());
- 完成后自动回到第 1 步,重新注册——形成可复用的工作循环。
⚠️ 注意:Worker.Start() 内部启动 goroutine 是合理设计(否则会阻塞),但调用方式应为 go worker.Start()(而非 worker.Start() 同步调用)。原文中 StartDispatcher 循环内调用 worker.Start() 实际已隐含 go,但语义不够清晰;更推荐显式写出 go worker.Start() 提升可读性与控制力。
? 请求调度流程(关键路径)
当 HTTP 请求到达 handleRequest:
- 解析请求体 → 构建 Request{Buf: buf, Delay: Delay};
- 发送至 RequestQueue
- 调度 goroutine(StartDispatcher 中的匿名函数)从 RequestQueue 接收请求;
- 从 WorkerQueue 取出某 Worker 的输入通道(worker :=
- 将请求发送至该通道:worker
- 对应 Worker 的 goroutine 从 w.Request 接收该请求,执行延时与写文件操作。
这个设计巧妙规避了“向未启动的通道发送数据”的风险——因为每个 Worker 在启动时就已将其 Request 通道放入 WorkerQueue,调度器永远能取到一个可用通道。
? 改进建议与最佳实践
-
显式启动 Worker goroutine(提升可维护性):
for i := 0; i
-
优雅关闭支持(生产环境必需):
// 关闭所有 Worker for i := 0; i
-
错误处理增强(当前 writeToFile 无错误检查):
if err := writeToFile(request.Buf); err != nil { log.Printf("worker%d: write failed: %v", w.ID, err) } 避免内存泄漏:ioutil.ReadAll(r.Body) 已被 io.ReadAll 替代(Go 1.16+),且需确保 r.Body.Close() 被调用(当前缺失,应补充)。
? 总结
该调度模型本质是 “通道的通道”(channel of channels) 模式:WorkerQueue 作为“工人资源池”,动态分配空闲 Worker;RequestQueue 作为“任务缓冲池”,解耦请求接收与处理速度差异。它轻量、可控、符合 Go 并发哲学——不要通过共享内存来通信,而应通过通信来共享内存。掌握此模式,即可快速构建高可靠、可伸缩的后台任务系统。











