如何在Golang gRPC中为不同接口设置客户端并发请求限制
Go gRPC服务实现(含接口并发限制)
核心思路
针对不同的gRPC接口,使用**信号量(Semaphore)**实现并发请求限制:
- 图片上传接口:用权重为10的信号量,限制同时最多处理10个请求
- 图片列表查询接口:用权重为100的信号量,限制同时最多处理100个请求
通过gRPC的**一元拦截器(Unary Interceptor)**将并发控制逻辑与业务逻辑解耦,无需在每个接口里重复写限制代码。
步骤1:定义gRPC服务Proto
编写image_service.proto文件,定义服务接口和消息结构:
syntax = "proto3"; package imageservice; option go_package = "./pb"; service ImageService { // 上传图片 rpc UploadImage(UploadImageRequest) returns (UploadImageResponse); // 查询图片列表 rpc ListImages(ListImagesRequest) returns (ListImagesResponse); } message UploadImageRequest { string filename = 1; bytes image_data = 2; } message UploadImageResponse { bool success = 1; string message = 2; } message ListImagesRequest {} message ListImagesResponse { repeated string filenames = 1; }
执行命令生成Go代码:
protoc --go_out=. --go-grpc_out=. image_service.proto
步骤2:实现并发限制拦截器
使用Go标准库sync/semaphore实现信号量,编写统一的拦截器根据方法名匹配对应的信号量:
package main import ( "context" "google.golang.org/grpc" "google.golang.org/grpc/codes" "google.golang.org/grpc/status" "sync/semaphore" ) // 定义两个信号量,分别对应不同接口的并发限制 var ( uploadSem = semaphore.NewWeighted(10) // 上传接口最多10并发 listSem = semaphore.NewWeighted(100) // 查询接口最多100并发 ) // concurrentLimitInterceptor 并发限制拦截器 func concurrentLimitInterceptor(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) { var sem *semaphore.Weighted // 根据gRPC方法名选择对应的信号量 switch info.FullMethod { case "/imageservice.ImageService/UploadImage": sem = uploadSem case "/imageservice.ImageService/ListImages": sem = listSem default: // 未知方法直接放行 return handler(ctx, req) } // 尝试获取信号量,无可用资源则返回限流错误 if !sem.TryAcquire(1) { return nil, status.Errorf(codes.ResourceExhausted, "请求过于频繁,请稍后再试") } defer sem.Release(1) // 请求处理完成后释放信号量 // 执行业务逻辑 return handler(ctx, req) }
步骤3:实现gRPC服务业务逻辑
编写服务的具体实现,处理图片上传和列表查询:
package main import ( "context" "io/ioutil" "os" "path/filepath" "your-project-path/pb" // 替换为你的proto生成的Go代码路径 ) const imageSaveDir = "./uploaded_images" type imageServiceServer struct { pb.UnimplementedImageServiceServer } // UploadImage 处理图片上传 func (s *imageServiceServer) UploadImage(ctx context.Context, req *pb.UploadImageRequest) (*pb.UploadImageResponse, error) { // 确保保存目录存在 if err := os.MkdirAll(imageSaveDir, 0755); err != nil { return &pb.UploadImageResponse{Success: false, Message: "创建目录失败"}, err } // 拼接保存路径 savePath := filepath.Join(imageSaveDir, req.Filename) // 写入图片文件 if err := ioutil.WriteFile(savePath, req.ImageData, 0644); err != nil { return &pb.UploadImageResponse{Success: false, Message: "保存图片失败"}, err } return &pb.UploadImageResponse{Success: true, Message: "图片上传成功"}, nil } // ListImages 返回已保存的图片列表 func (s *imageServiceServer) ListImages(ctx context.Context, req *pb.ListImagesRequest) (*pb.ListImagesResponse, error) { // 读取目录下的所有文件 files, err := ioutil.ReadDir(imageSaveDir) if err != nil { return &pb.ListImagesResponse{}, err } var filenames []string for _, file := range files { if !file.IsDir() { filenames = append(filenames, file.Name()) } } return &pb.ListImagesResponse{Filenames: filenames}, nil }
步骤4:启动gRPC服务
在main函数中注册服务并添加拦截器:
package main import ( "net" "google.golang.org/grpc" "your-project-path/pb" ) func main() { // 监听端口 lis, err := net.Listen("tcp", ":50051") if err != nil { panic("failed to listen: " + err.Error()) } // 创建gRPC服务器,添加并发限制拦截器 grpcServer := grpc.NewServer( grpc.UnaryInterceptor(concurrentLimitInterceptor), ) // 注册服务 pb.RegisterImageServiceServer(grpcServer, &imageServiceServer{}) // 启动服务 if err := grpcServer.Serve(lis); err != nil { panic("failed to serve: " + err.Error()) } }
关键注意事项
- 信号量释放必须用
defer,确保即使业务逻辑出错也能释放资源,避免死锁 - 若需要等待超时逻辑,可替换
TryAcquire为Acquire:if err := sem.Acquire(ctx, 1); err != nil { return nil, status.Errorf(codes.Canceled, "请求被取消或超时") } - 图片保存目录建议配置为环境变量,方便部署时调整
- 生产环境中建议添加日志,记录并发请求的获取和释放情况,便于排查问题
内容的提问来源于stack exchange,提问作者Islom
相关产品推荐
相关产品推荐

