You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在两个限频独立端点间同步请求?

多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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 11:03:08