Go-Zero项目开发36: 实现社交服务创群请求幂等性
纲要
- 幂等性的意义与场景
- 基于唯一请求 ID 的幂等方案
- 核心调用流程(搭配时序图)
- 项目代码结构
- 幂等组件实现
- 幂等接口定义
Idempotent - 默认实现(基于 Redis 与缓存)
- RPC 客户端拦截器
- RPC 服务端拦截器
- HTTP API 中间件
- 幂等接口定义
- 集成到社交服务(创建群组)
- API 层配置
- RPC 服务端配置
- RPC 客户端配置
- 测试与效果验证
- 总结
幂等性的意义与场景
在微服务架构中,一次业务操作可能会因为网络波动、服务临时不可用等原因发生超时,调用方往往会发起重试。如果被调用的服务没有做幂等保护,同样的请求可能被处理多次,例如创建群组时出现两个相同的群,这类“写”操作对数据一致性的破坏是致命的。
幂等性(Idempotence)就是用来解决这一问题的:同一个操作执行一次或多次,产生的业务效果和数据影响完全一致,不会因为重复调用而产生副作用。
本章在 go-zero 社交服务中,为“创建群组”接口引入幂等性,确保即便上游因为超时重试,最终也只会真正创建一个群。
基于唯一请求 ID 的幂等方案
核心思路是为每一次业务请求分配一个全局唯一的请求 ID,服务端在处理前先检查该 ID 是否已经处理过:
- 若 ID 不存在,则执行正常业务逻辑,并将执行结果与请求 ID 绑定存储;
- 若 ID 已存在,说明该请求正在处理或已经完成,直接返回之前的结果或给出友好提示,不再重复执行业务。
请求 ID 的生成与传递过程如下:
- 客户端(API 网关或上游服务)在发起请求时,通过HTTP 中间件生成唯一请求 ID,并注入到
context中; - 当通过 gRPC 调用下游服务时,RPC 客户端拦截器从
context取出请求 ID,写入到 gRPC 的metadata中; - RPC 服务端拦截器从
metadata中取出请求 ID,调用幂等组件完成“是否已处理”的判断与结果保存。
整个流程可以由下面的时序图概括:
项目代码结构
我们新增一个公共包pkg/idempotent来实现幂等组件,同时定义 HTTP 中间件和 RPC 拦截器。文件组织如下:
pkg/ └─ idempotent/ ├─ idempotent.go # 接口定义与默认实现 ├─ clientinterceptor.go # gRPC 客户端拦截器 └─ serverinterceptor.go # gRPC 服务端拦截器 api/ └─ internal/ └─ middleware/ └─ requestid.go # HTTP 请求 ID 中间件幂等组件实现
幂等接口定义Idempotent
首先定义幂等处理的核心能力:获取请求标识、判断方法是否支持幂等、校验幂等状态、保存结果。
packageidempotentimport("context")// Idempotent 定义幂等性处理接口typeIdempotentinterface{// GetReqId 从上下文中获取请求的唯一标识GetReqId(ctx context.Context,methodstring)string// SupportIdempotent 判断指定方法是否需要幂等保护SupportIdempotent(methodstring)bool// CheckIdempotent 校验幂等:如果任务尚未执行,返回 false;否则返回已有结果或错误CheckIdempotent(ctx context.Context,idstring,methodstring)(executingbool,resultinterface{},errerror)// SaveResult 保存执行结果,与 id 绑定SaveResult(ctx context.Context,idstring,methodstring,resultinterface{},errerror)error}默认实现(基于 Redis 与缓存)
默认实现利用 Redis 的SETNX作为分布式锁,保证同一个请求 ID 只被一个处理流程接受;同时使用 go-zero 的Cache存储执行结果,并设置过期时间,避免数据无限堆积。
packageidempotentimport("context""fmt""time""github.com/zeromicro/go-zero/core/collection""github.com/zeromicro/go-zero/core/stores/cache""github.com/zeromicro/go-zero/core/stores/redis")// 默认幂等处理对象typedefaultIdempotentstruct{rds*redis.Redis cache*cache.Cache idempotentMethodsmap[string]bool}// NewIdempotent 创建一个默认幂等实例funcNewIdempotent(redisConf redis.RedisConf,ttl time.Duration,methods...string)Idempotent{di:=&defaultIdempotent{rds:redis.MustNewRedis(redisConf),idempotentMethods:make(map[string]bool),}for_,m:=rangemethods{di.idempotentMethods[m]=true}// 使用 go-zero 缓存,设置默认过期时间c:=cache.New(cache.WithExpiry(ttl))di.cache=creturndi}// GetReqId 从 context 中获取请求 ID,若不存在则基于方法名和随机串生成func(d*defaultIdempotent)GetReqId(ctx context.Context,methodstring)string{// 尝试从 context 取值,若已有则直接返回(通常在中间件中已经设置)ifid,ok:=ctx.Value("requestId").(string);ok&&id!=""{returnid}// 兜底生成:method + 随机数,实际场景可使用 UUIDreturnfmt.Sprintf("%s-%d",method,time.Now().UnixNano())}// SupportIdempotent 判断当前方法是否需要幂等处理func(d*defaultIdempotent)SupportIdempotent(methodstring)bool{returnd.idempotentMethods[method]}// CheckIdempotent 使用 Redis SETNX 实现幂等校验func(d*defaultIdempotent)CheckIdempotent(ctx context.Context,idstring,methodstring)(bool,interface{},error){key:=fmt.Sprintf("idempotent:%s:%s",method,id)// 尝试设置 key,过期时间防止死锁ok,err:=d.rds.SetnxEx(key,"1",10*time.Second)iferr!=nil{returnfalse,nil,err}ifok{// 设置成功,说明是首次处理returnfalse,nil,nil}// key 已存在,尝试从缓存获取结果varresultinterface{}err=d.cache.Get(key,&result)iferr!=nil{// 缓存中没有结果,说明正在执行中returntrue,nil,fmt.Errorf("任务正在执行中,请稍后重试")}// 已有结果,直接返回returntrue,result,nil}// SaveResult 保存执行结果到缓存func(d*defaultIdempotent)SaveResult(ctx context.Context,idstring,methodstring,resultinterface{},execErrerror)error{key:=fmt.Sprintf("idempotent:%s:%s",method,id)returnd.cache.SetWithExpire(key,result,10*time.Minute)}RPC 客户端拦截器
客户端拦截器的职责是在发起 gRPC 调用前,从上下文中取出请求 ID,并将其写入到 outgoing metadata 中。
packageidempotentimport("context""google.golang.org/grpc""google.golang.org/grpc/metadata")// ClientInterceptor 返回一个 gRPC 客户端拦截器funcClientInterceptor(idempotent Idempotent)grpc.UnaryClientInterceptor{returnfunc(ctx context.Context,methodstring,req,replyinterface{},cc*grpc.ClientConn,invoker grpc.UnaryInvoker,opts...grpc.CallOption)error{id:=idempotent.GetReqId(ctx,method)// 将请求 ID 注入到 gRPC metadata 中md:=metadata.Pairs("x-request-id",id)ctx=metadata.NewOutgoingContext(ctx,md)returninvoker(ctx,method,req,reply,cc,opts...)}}RPC 服务端拦截器
服务端拦截器从 metadata 取出请求 ID,调用幂等组件判断:
- 未执行过:放行到业务逻辑,并在业务执行后将结果通过
SaveResult保存; - 正在执行中:返回错误提示;
- 已完成:直接返回已保存的结果。
packageidempotentimport("context""fmt""google.golang.org/grpc""google.golang.org/grpc/metadata""google.golang.org/grpc/status""google.golang.org/grpc/codes")// ServerInterceptor 返回一个 gRPC 服务端拦截器funcServerInterceptor(idempotent Idempotent)grpc.UnaryServerInterceptor{returnfunc(ctx context.Context,reqinterface{},info*grpc.UnaryServerInfo,handler grpc.UnaryHandler)(interface{},error){method:=info.FullMethod// 从 metadata 中提取请求 IDmd,ok:=metadata.FromIncomingContext(ctx)if!ok{returnnil,status.Error(codes.InvalidArgument,"缺少 metadata")}ids:=md.Get("x-request-id")iflen(ids)==0{returnnil,status.Error(codes.InvalidArgument,"缺少请求 ID")}id:=ids[0]// 若该方法无需幂等保护,直接执行if!idempotent.SupportIdempotent(method){returnhandler(ctx,req)}// 幂等校验executing,result,err:=idempotent.CheckIdempotent(ctx,id,method)iferr!=nil{returnnil,status.Errorf(codes.ResourceExhausted,"幂等校验失败: %v",err)}ifexecuting{ifresult!=nil{// 之前已执行完,返回缓存结果returnresult,nil}// 正在执行中returnnil,status.Error(codes.Aborted,"任务正在执行中,请稍后重试")}// 正常执行业务resp,err:=handler(ctx,req)// 保存结果(无论成功失败都保存,以便重试时快速返回)ifsaveErr:=idempotent.SaveResult(ctx,id,method,resp,err);saveErr!=nil{fmt.Printf("保存幂等结果失败: %v\n",saveErr)}returnresp,err}}HTTP API 中间件
API 网关作为 HTTP 请求的入口,通过中间件为每个请求生成唯一 ID,并存入context,后续所有 RPC 调用都能通过客户端拦截器携带此 ID。
packagemiddlewareimport("context""net/http""github.com/google/uuid")// RequestIdMiddleware 为每个 HTTP 请求生成唯一的请求 ID 并注入 contextfuncRequestIdMiddleware(next http.HandlerFunc)http.HandlerFunc{returnfunc(w http.ResponseWriter,r*http.Request){// 优先使用客户端传递的 X-Request-Id,否则生成新 IDreqId:=r.Header.Get("X-Request-Id")ifreqId==""{reqId=uuid.New().String()}ctx:=context.WithValue(r.Context(),"requestId",reqId)next(w,r.WithContext(ctx))}}集成到社交服务(创建群组)
API 层配置
在创建群组的 API 路由上应用上述中间件,同时确保 RPC 客户端配置了幂等拦截器。
// api/internal/handler/group/create_group_handler.go 示例片段funcRegisterHandlers(server*rest.Server,ctx*svc.ServiceContext){server.AddRoutes([]rest.Route{{Method:http.MethodPost,Path:"/group/create",Handler:middleware.RequestIdMiddleware(CreateGroupHandler(ctx)),},},)}对应的CreateGroupHandler内部会调用 RPC 客户端,该客户端在初始化时已经添加了客户端拦截器。
RPC 客户端配置
在ServiceContext中创建 RPC 客户端时,通过zrpc.MustNewClient的选项添加拦截器。
// api/internal/svc/service_context.go 示例import("github.com/zeromicro/go-zero/zrpc""your_project/pkg/idempotent")typeServiceContextstruct{Config config.Config GroupRpc group.GroupClient// 生成的 gRPC 客户端}funcNewServiceContext(c config.Config)*ServiceContext{// 初始化幂等组件(指定需要幂等保护的方法)idemp:=idempotent.NewIdempotent(c.Redis,10*time.Minute,"/group.Group/CreateGroup")// 创建 gRPC 客户端,并注入拦截器conn:=zrpc.MustNewClient(c.GroupRpc,zrpc.WithUnaryClientInterceptor(idempotent.ClientInterceptor(idemp)))return&ServiceContext{Config:c,GroupRpc:group.NewGroupClient(conn),}}RPC 服务端配置
在 RPC 服务的启动文件中,为 gRPC 服务端添加服务端拦截器。
// rpc/internal/server/group_server.go 或 main.gofuncmain(){varc config.Config conf.MustLoad("etc/group.yaml",&c)idemp:=idempotent.NewIdempotent(c.Redis,10*time.Minute,"/group.Group/CreateGroup")s:=zrpc.MustNewServer(c.RpcServerConf,func(grpcServer*grpc.Server){group.RegisterGroupServer(grpcServer,server.NewGroupServer(svc.NewServiceContext(c)))},zrpc.WithUnaryServerInterceptor(idempotent.ServerInterceptor(idemp)))defers.Stop()s.Start()}测试与效果验证
启动 API 服务和 RPC 服务后,使用工具连续发送两次完全相同的创建群组请求(相同的请求 ID):
- 第一次请求:服务端日志显示进入幂等校验,SETNX 成功,执行创建群组业务,数据库中新增一条记录,同时结果被缓存;
- 第二次请求(相同的请求 ID):服务端日志显示
任务已存在,直接从缓存返回第一次的结果,不会再次操作数据库。
数据库中也只会保留一条群组记录,证明幂等性生效。如果两次请求间隔极短,第二次请求可能捕获到“正在执行中”的状态,返回提示,避免重复写入。
总结
在 go-zero 框架下,通过“唯一请求 ID + Redis SETNX + 拦截器”的方式,为社交服务的“创建群组”等关键写操作实现了可靠的幂等性。幂等组件被封装为可复用的公共模块,仅需通过配置文件指定需要保护的方法,即可无缝集成到任何 gRPC 服务中。
该方案同时兼顾了:
- 防止重复执行:利用 Redis 原子操作;
- 结果快速返回:缓存已完成的结果,避免业务重新计算;
- 解耦与可扩展:通过 context 和 metadata 传递 ID,不侵入业务代码。
