Skip to content

Day 3:请求队列+goroutine池——削峰填谷

今天做什么

昨天限流是"拒绝",今天加请求队列——把超出限流的请求排队,而不是直接拒绝。再加上 goroutine 池控制并发。

代码

go
package main

import (
    "context"
    "sync"
)

// ===== 请求队列 =====
type RequestQueue struct {
    queue    chan *QueuedRequest
    maxSize  int
    dropped  atomic.Int64
}

type QueuedRequest struct {
    Request  *http.Request
    Writer   http.ResponseWriter
    Done     chan struct{}
}

func NewRequestQueue(maxSize int) *RequestQueue {
    return &RequestQueue{
        queue:   make(chan *QueuedRequest, maxSize),
        maxSize: maxSize,
    }
}

func (q *RequestQueue) Enqueue(w http.ResponseWriter, 
    r *http.Request) bool {
    select {
    case q.queue <- &QueuedRequest{
        Request: r,
        Writer:  w,
        Done:    make(chan struct{}),
    }:
        return true
    default:
        q.dropped.Add(1)
        return false // 队列满了,丢弃
    }
}

// ===== Goroutine 池 =====
type WorkerPool struct {
    workers int
    handler func(*QueuedRequest)
}

func NewWorkerPool(workers int, handler func(*QueuedRequest)) *WorkerPool {
    return &WorkerPool{
        workers: workers,
        handler: handler,
    }
}

func (p *WorkerPool) Start(ctx context.Context) {
    var wg sync.WaitGroup
    for i := 0; i < p.workers; i++ {
        wg.Add(1)
        go func(id int) {
            defer wg.Done()
            for {
                select {
                case <-ctx.Done():
                    return
                case task := <-taskQueue:
                    p.handler(task)
                }
            }
        }(i)
    }
}

// ===== 网关集成 =====
func (g *Gateway) ServeHTTP(w http.ResponseWriter, r *http.Request) {
    // 第一关:限流
    if !g.limiter.Allow() {
        // 不限流就不直接拒绝,而是排队
        if !g.queue.Enqueue(w, r) {
            http.Error(w, `{"error":"overloaded"}`, 503)
        }
        return
    }

    // 第二关:熔断
    if !g.breaker.Allow() {
        http.Error(w, `{"error":"circuit open"}`, 503)
        return
    }

    g.router.ProxyHandler(w, r)
}

func main() {
    gw := &Gateway{
        router:  NewRouter(),
        limiter: NewTokenBucket(100, 200),
        breaker: NewCircuitBreaker(10, 30*time.Second),
        queue:   NewRequestQueue(500), // 最多排队500个
    }

    // 启动 Worker 池处理排队的请求
    pool := NewWorkerPool(20, gw.processQueuedRequest)
    go pool.Start(context.Background())

    http.ListenAndServe(":8080", http.HandlerFunc(gw.ServeHTTP))
}

效果

限流前:每秒100个以内 → 正常处理
限流后:超出的 → 排队等待(最多500个)
队列满 → 503 拒绝

Day 3 完成。网关会削峰填谷了。