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调用幂等组件完成“是否已处理”的判断与结果保存。整个流程可以由下面的时序图概括业务幂等组件(Redis)RPC 服务端拦截器RPC 客户端拦截器HTTP 中间件客户端/API业务幂等组件(Redis)RPC 服务端拦截器RPC 客户端拦截器HTTP 中间件客户端/APIalt[未执行过][已执行或正在执行]请求创建群生成唯一请求ID注入 ctx携带 ctx 发起 RPC 调用从 ctx 获取请求ID通过 metadata 传递请求IDCheckIdempotent(ctx, id, method)false, nil执行创建群逻辑返回结果SaveResult(ctx, id, method, result)true / 已有结果返回已有结果或提示响应响应响应项目代码结构我们新增一个公共包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(contextfmttimegithub.com/zeromicro/go-zero/core/collectiongithub.com/zeromicro/go-zero/core/stores/cachegithub.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.cachecreturndi}// GetReqId 从 context 中获取请求 ID若不存在则基于方法名和随机串生成func(d*defaultIdempotent)GetReqId(ctx context.Context,methodstring)string{// 尝试从 context 取值若已有则直接返回通常在中间件中已经设置ifid,ok:ctx.Value(requestId).(string);okid!{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{}errd.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(contextgoogle.golang.org/grpcgoogle.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)ctxmetadata.NewOutgoingContext(ctx,md)returninvoker(ctx,method,req,reply,cc,opts...)}}RPC 服务端拦截器服务端拦截器从 metadata 取出请求 ID调用幂等组件判断未执行过放行到业务逻辑并在业务执行后将结果通过SaveResult保存正在执行中返回错误提示已完成直接返回已保存的结果。packageidempotentimport(contextfmtgoogle.golang.org/grpcgoogle.golang.org/grpc/metadatagoogle.golang.org/grpc/statusgoogle.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(contextnet/httpgithub.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{reqIduuid.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/zrpcyour_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)))returnServiceContext{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不侵入业务代码。