如何在两个限频独立端点间同步请求?
多API速率与并发控制解决方案
需求概述
- Endpoint 1(GET):速率限制10次/秒,每次返回2~3000个对象的数组
- Endpoint 2(POST):速率限制20次/秒,需接收Endpoint1返回的每个对象的部分数据
- 需控制每个端点的并发请求数,同时为500+等失败请求添加重试机制
完整实现代码
package main import ( "bytes" "encoding/json" "fmt" "io" "net/http" "sync" "time" ) type HttpBinGetRequest struct { url string } type HttpBinGetResponse struct { Uuid string `json:"uuid"` StatusCode int } type HttpBinPostRequest struct { url string uuid string // Item to post to API } type HttpBinPostResponse struct { Data string `json:"data"` StatusCode int } func main() { // Prepare GET requests for 500 requests var requests []*HttpBinGetRequest for i := 0; i < 500; i++ { uri := "https://httpbin.org/uuid" requests = append(requests, &HttpBinGetRequest{url: uri}) } // -------------------------- // GET 请求控制:速率10次/秒,并发数10 // -------------------------- getSemaphore := make(chan struct{}, 10) getRate := make(chan struct{}, 10) for i := 0; i < cap(getRate); i++ { getRate <- struct{}{} } // 速率补充定时器:每100ms放回一个令牌,保证每秒最多10次请求 go func() { ticker := time.NewTicker(100 * time.Millisecond) defer ticker.Stop() for range ticker.C { select { case getRate <- struct{}{}: default: } } }() var wg sync.WaitGroup // 收集GET成功返回的UUID var uuids []string // 用于并发安全地写入uuids var uuidMu sync.Mutex for _, request := range requests { wg.Add(1) go func(r *HttpBinGetRequest) { defer wg.Done() // 等待速率令牌 <-getRate // 等待并发令牌 getSemaphore <- struct{}{} defer func() { <-getSemaphore }() // 带重试的GET请求 resp, err := retry(func() (*HttpBinGetResponse, error) { return get(r) }, 3) if err != nil { fmt.Printf("GET请求失败:%v\n", err) return } if resp.StatusCode == http.StatusOK { uuidMu.Lock() uuids = append(uuids, resp.Uuid) uuidMu.Unlock() fmt.Printf("GET成功获取UUID:%s\n", resp.Uuid) } else { fmt.Printf("GET请求返回非200状态:%d\n", resp.StatusCode) } }(request) } wg.Wait() fmt.Printf("共获取到%d个有效UUID\n", len(uuids)) // -------------------------- // POST 请求控制:速率20次/秒,并发数20 // -------------------------- postSemaphore := make(chan struct{}, 20) postRate := make(chan struct{}, 20) for i := 0; i < cap(postRate); i++ { postRate <- struct{}{} } // 速率补充定时器:每50ms放回一个令牌,保证每秒最多20次请求 go func() { ticker := time.NewTicker(50 * time.Millisecond) defer ticker.Stop() for range ticker.C { select { case postRate <- struct{}{}: default: } } }() var postWg sync.WaitGroup for _, uuid := range uuids { postWg.Add(1) go func(u string) { defer postWg.Done() // 等待速率令牌 <-postRate // 等待并发令牌 postSemaphore <- struct{}{} defer func() { <-postSemaphore }() req := &HttpBinPostRequest{ url: "https://httpbin.org/post", uuid: u, } // 带重试的POST请求 resp, err := retry(func() (*HttpBinPostResponse, error) { return post(req) }, 3) if err != nil { fmt.Printf("POST UUID %s失败:%v\n", u, err) return } if resp.StatusCode == http.StatusOK { fmt.Printf("POST UUID %s成功,响应数据:%s\n", u, resp.Data) } else { fmt.Printf("POST UUID %s返回非200状态:%d\n", u, resp.StatusCode) } }(uuid) } postWg.Wait() fmt.Println("所有POST请求处理完成") } // 通用重试函数:针对失败请求进行指定次数的重试 func retry[T any](fn func() (T, error), maxRetries int) (T, error) { var result T for i := 0; i < maxRetries; i++ { var err error result, err = fn() if err == nil { return result, nil } // 重试间隔,指数退避 time.Sleep(time.Duration(i+1) * 500 * time.Millisecond) fmt.Printf("请求失败,正在进行第%d次重试...\n", i+1) } return result, fmt.Errorf("已达到最大重试次数%d,请求仍失败", maxRetries) } func get(hbgr *HttpBinGetRequest) (*HttpBinGetResponse, error) { httpResp := &HttpBinGetResponse{} client := &http.Client{Timeout: 5 * time.Second} req, err := http.NewRequest("GET", hbgr.url, nil) if err != nil { return httpResp, fmt.Errorf("创建GET请求失败:%w", err) } req.Header.Set("accept", "application/json") resp, err := client.Do(req) if err != nil { return httpResp, fmt.Errorf("发送GET请求失败:%w", err) } defer resp.Body.Close() if resp.StatusCode >= 500 { return httpResp, fmt.Errorf("服务器错误,状态码:%d", resp.StatusCode) } body, err := io.ReadAll(resp.Body) if err != nil { return httpResp, fmt.Errorf("读取GET响应体失败:%w", err) } if err := json.Unmarshal(body, &httpResp); err != nil { return httpResp, fmt.Errorf("解析GET响应失败:%w", err) } httpResp.StatusCode = resp.StatusCode return httpResp, nil } func post(hbr *HttpBinPostRequest) (*HttpBinPostResponse, error) { httpResp := &HttpBinPostResponse{} client := &http.Client{Timeout: 5 * time.Second} req, err := http.NewRequest("POST", hbr.url, bytes.NewBuffer([]byte(hbr.uuid))) if err != nil { return httpResp, fmt.Errorf("创建POST请求失败:%w", err) } req.Header.Set("accept", "application/json") req.Header.Set("Content-Type", "text/plain") resp, err := client.Do(req) if err != nil { return httpResp, fmt.Errorf("发送POST请求失败:%w", err) } defer resp.Body.Close() if resp.StatusCode == http.StatusTooManyRequests { retryAfter := resp.Header.Get("Retry-After") if retryAfter != "" { fmt.Printf("触发速率限制,需等待%s秒\n", retryAfter) } return httpResp, fmt.Errorf("速率限制触发,状态码:%d", resp.StatusCode) } if resp.StatusCode >= 500 { return httpResp, fmt.Errorf("服务器错误,状态码:%d", resp.StatusCode) } body, err := io.ReadAll(resp.Body) if err != nil { return httpResp, fmt.Errorf("读取POST响应体失败:%w", err) } if err := json.Unmarshal(body, &httpResp); err != nil { return httpResp, fmt.Errorf("解析POST响应失败:%w", err) } httpResp.StatusCode = resp.StatusCode return httpResp, nil }
关键改进说明
- 并发与速率分离控制:
- 每个端点使用两个通道:
Semaphore控制并发请求数(避免同时发起过多请求压垮自身或服务器),Rate通道配合定时器实现精确的每秒请求次数限制 - GET端点每100ms补充一个速率令牌,保证每秒最多10次;POST端点每50ms补充一个,保证每秒最多20次
- 每个端点使用两个通道:
- 并发安全收集数据:使用
sync.Mutex保证多个协程写入UUID切片时的线程安全 - 通用重试机制:封装
retry函数,对500+服务器错误等失败请求进行最多3次重试,采用指数退避间隔避免加重服务器负担 - 错误处理增强:在
get和post函数中细化错误类型,区分请求创建失败、发送失败、服务器错误、解析错误等场景,便于排查问题 - 超时设置:为HTTP客户端添加5秒超时,避免请求长时间阻塞
内容的提问来源于stack exchange,提问作者Coldchain9
相关产品推荐
相关产品推荐

