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

如何在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())
  }
}

关键注意事项

  1. 信号量释放必须用defer,确保即使业务逻辑出错也能释放资源,避免死锁
  2. 若需要等待超时逻辑,可替换TryAcquire为Acquire:
    if err := sem.Acquire(ctx, 1); err != nil {
      return nil, status.Errorf(codes.Canceled, "请求被取消或超时")
    }
    
  3. 图片保存目录建议配置为环境变量,方便部署时调整
  4. 生产环境中建议添加日志,记录并发请求的获取和释放情况,便于排查问题

内容的提问来源于stack exchange,提问作者Islom

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 10:38:34