GoFr 如何用 gofr wrap grpc 实现 gRPC 双向流式聊天服务
GoFr 如何用 gofr wrap grpc 实现 gRPC 双向流式聊天服务【免费下载链接】gofrAn opinionated GoLang framework for accelerated microservice development. Built in support for databases and observability.项目地址: https://gitcode.com/GitHub_Trending/go/gofr如果你的目标是实现一个聊天服务让客户端和服务端在同一个 RPC 上互相发送多条消息实时聊天、交互式协议GoFr 提供的落地路径是在 proto 中用stream关键字声明双向流方法BiDiStream再用gofr wrap grpc生成带 GoFr 上下文的服务端与客户端封装你只写业务逻辑tracing、metrics 和日志由生成的 wrapper 自动接入。本文以仓库中的 grpc-streaming-server 与 grpc-streaming-client 两个示例工程为主线走通从 proto 定义到双向流验证的完整流程。完整说明见 gRPC Streaming 文档gofr wrap grpc的参数说明见 CLI 参考。准备条件按 gRPC Streaming 文档 的 Prerequisites 一节实现 gRPC 流式服务前需要安装 Protocol Buffer 编译器protoc3 版本安装 Go gRPC 代码生成插件go install google.golang.org/protobuf/cmd/protoc-gen-gov1.28 go install google.golang.org/grpc/cmd/protoc-gen-go-grpcv1.2 export PATH$PATH:$(go env GOPATH)/bin安装 gofr-cligo install gofr.dev/cli/gofrlatest如果你直接运行仓库自带示例chat.pb.go、chat_grpc.pb.go等 protoc 生成文件已存在于 server 目录可以跳过 protoc 安装直接从生成命令开始。定义带 BiDiStream 的 proto 文件双向流的声明方式是在.proto文件里给入参和出参加上stream关键字。仓库示例 chat.proto 定义了ChatService包含三种流式 RPC其中BiDiStream就是本文目标的双向流方法syntax proto3; option go_package gofr.dev/examples/grpc/grpc-streaming-server/server; message Request { string message 1; } message Response { string message 1; } service ChatService { rpc ServerStream(Request) returns (stream Response); rpc ClientStream(stream Request) returns (Response); rpc BiDiStream(stream Request) returns (stream Response); }go_package指向存放生成代码的 Go 包路径文档示例中写的是path/to/your/proto/file在自己的项目里替换为自己的包路径。示例保留了三种流式方法本文只聚焦BiDiStream的实现与验证。用 gofr wrap grpc 生成服务端与客户端封装在服务端工程目录执行gofr wrap grpc server -proto./path/to/your/proto/file命令以 proto 文件为输入为本文的ChatService生成chatservice_server.go待填业务逻辑的模板文件示例中的 chatservice_server.go 即生成后填充了实现chatservice_gofr.go自动生成的 wrapper文件头标注Code generated by gofr.dev/cli/gofr. DO NOT EDIT.包含流式方法的埋点包装不要手改request_gofr.go请求包装用于绑定上下文health_gofr.go健康检查服务集成。在客户端工程执行对应命令生成客户端封装gofr wrap grpc client -proto./path/to/your/proto/file生成chatservice_client.go。示例客户端的生成文件 定义了ChatServiceGoFrClient接口双向流对应的方法签名是BiDiStream(ctx *gofr.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[Request, Response], error)实现双向流的服务端逻辑chatservice_server.go里的BiDiStream方法是你要填的核心。示例实现完整代码见 chatservice_server.go是一个回声服务在 goroutine 中循环Recv()收消息收到什么回什么func (s *ChatServiceGoFrServer) BiDiStream(ctx *gofr.Context, stream ChatService_BiDiStreamServer) error { // Handle incoming messages in a goroutine errChan : make(chan error) go func() { for { // Check if context is canceled select { case -stream.Context().Done(): errChan - status.Error(codes.Canceled, client disconnected) return default: } req, err : stream.Recv() if err io.EOF { break } if err ! nil { errChan - status.Errorf(codes.Internal, error receiving stream: %v, err) return } // Process request and send response resp : Response{Message: Echo: req.Message} if err : stream.Send(resp); err ! nil { errChan - status.Errorf(codes.Internal, error sending stream: %v, err) return } } errChan - nil }() // Wait for completion select { case err : -errChan: return err case -stream.Context().Done(): return status.Error(codes.Canceled, client disconnected) } }文档给出的双向流实现要点对照这段代码理解用 goroutine 处理收消息外层select等待完成或取消每次Recv()前检查stream.Context().Done()以检测客户端断开Recv()返回io.EOF表示客户端结束发送是流的正常结束方式用break退出循环出错时返回带 gRPC status code 的错误取消用codes.Canceled收发异常用codes.Internal用errChan在 goroutine 与主逻辑之间协调结果。在 main.go 注册服务并启动示例的 main.go 只有几行流式服务的注册方式和 unary 服务一样package main import ( gofr.dev/examples/grpc/grpc-streaming-server/server gofr.dev/pkg/gofr ) func main() { app : gofr.New() // Register streaming service server.RegisterChatServiceServerWithGofr(app, server.NewChatServiceGoFrServer()) app.Run() }RegisterChatServiceServerWithGofr是chatservice_gofr.go生成的注册函数内部会把你的实现包进带埋点的 wrapper 再注册到 gRPC server。服务端运行配置在 configs/.envgRPC 服务监听GRPC_PORT9000APP_NAMEgrpc-server-example GRPC_PORT9000 HTTP_PORT8081 LOG_LEVELDEBUG TRACE_EXPORTERgofr在examples/grpc/grpc-streaming-server目录执行go run main.go启动示例 README 给出的就是这条命令。编写客户端并发起双向流调用客户端示例grpc-streaming-client/main.go用生成的客户端创建 ChatService 连接并把三种流式调用分别暴露成 HTTP 端点双向流对应GET /chat/bidi-streamfunc main() { app : gofr.New() // Create a gRPC client for the Chat Streaming service chatClient, err : client.NewChatServiceGoFrClient(app.Config.Get(GRPC_SERVER_HOST), app.Metrics()) if err ! nil { app.Logger().Errorf(Failed to create Chat client: %v, err) } chat : NewChatHandler(chatClient) app.GET(/chat/server-stream, chat.ServerStreamHandler) app.POST(/chat/client-stream, chat.ClientStreamHandler) app.GET(/chat/bidi-stream, chat.BiDiStreamHandler) app.Run() }客户端配置 configs/.env 中GRPC_SERVER_HOSTlocalhost:9000指向服务端 gRPC 端口HTTP_PORT8080是客户端自身的 HTTP 端点端口。BiDiStreamHandler展示了双向流客户端的完整调用模式先BiDiStream(ctx)拿到流开 goroutine 循环stream.Recv()收响应主流程依次stream.Send()发消息发完调用stream.CloseSend()关闭发送侧再收集响应stream, err : c.chatClient.BiDiStream(ctx) if err ! nil { return nil, fmt.Errorf(failed to initiate bidirectional stream: %v, err) } respChan, errChan : make(chan StreamResponse), make(chan error) go c.receiveBiDiResponses(ctx, stream, respChan, errChan) sentMessages, err : c.sendBiDiMessages(ctx, stream, streamLog) if err ! nil { return nil, err } if err : stream.CloseSend(); err ! nil { return nil, fmt.Errorf(failed to close send: %v, err) } receivedMessages, err : c.collectBiDiResponses(respChan, errChan, streamLog)配套辅助函数同一文件中sendBiDiMessages依次发送message 1、message 2、message 3receiveBiDiResponses在 goroutine 里循环Recv()收到io.EOF时向errChan发nil表示流正常结束collectBiDiResponses从respChan收集响应带 5 秒超时保护超时返回bidirectional stream timeout错误。验证双向流服务验证方式一运行服务端测试。示例仓库自带 main_test.go其中TestBiDiStream直接向 gRPC 端口发起双向流调用发送msg1、msg2、msg3并CloseSend()然后逐条断言收到的响应为Echo: msg1、Echo: msg2、Echo: msg3。测试文件的TestMain会先以 goroutine 启动main()用testutil.ReserveServerPorts()预留空闲端口代码注释说明固定端口 9000 在 CI 中会与 MinIO 冲突系统环境变量优先于配置文件再用testutil.WaitForGRPCServerE等待服务就绪后运行测试。在服务端目录执行go test ./...TestBiDiStream通过即说明双向流服务端按预期工作。同文件还覆盖了 context 取消预期错误 status 为codes.Canceled、超时、EOF 处理等场景。验证方式二跑通客户端示例。按 客户端 README先在examples/grpc/grpc-streaming-server执行go run main.go启动服务端再在examples/grpc/grpc-streaming-client执行go run main.go启动客户端然后请求curl http://localhost:8080/chat/bidi-stream按 客户端 handler 代码响应 JSON 包含status值为bidirectional stream completed、sent_messages、received_messages、detailed_log等字段由于服务端是回声实现received_messages中应能看到Echo: message 1、Echo: message 2、Echo: message 3。客户端每收发一条消息还会打日志如Received bidirectional message: ...可配合终端输出核对消息流向。内置可观测性与流错误处理生成代码已为每个流式操作接入观测能力无需额外配置Metrics服务端注册app_gRPC-Stream_statsSend、Recv、SendAndClose、CloseSend 的耗时直方图客户端注册app_gRPC-Client-Stream_statsTracing每次流操作Send、Recv 等自动创建 span用于分布式追踪Logging流操作自动记录操作类型、方法名、耗时与错误状态。以 chatservice_gofr.go 为例bidiStreamWrapperBiDiStream的Send/Recv/CloseSend每个方法都先ctx.Trace(method /Send)创建 span再调用DocumentRPCLog写入app_gRPC-Stream_stats指标。流式错误处理上文档明确了两类情况io.EOF是流的正常结束双向流中客户端CloseSend()后服务端Recv()收到 EOF客户端侧Recv()收到 EOF 表示服务端结束发送context 取消通过stream.Context().Done()检测并返回合适的 status code仓库测试对取消场景的断言就是验证status.FromError(err).Code() codes.Canceled。另外GoFr 支持通过app.AddGRPCServerStreamInterceptors为所有流式 RPC 添加 stream interceptor用于贯穿整个流生命周期的逻辑如流级别鉴权用法见 gRPC Streaming 文档 的 Adding Custom Stream interceptors 一节。参考资料服务端示例examples/grpc/grpc-streaming-serverproto 定义见 chat.proto客户端示例examples/grpc/grpc-streaming-client生成命令说明gofr wrap grpc本文只覆盖了双向流BiDiStream。同一 proto 里的ServerStream服务端多次stream.Send()与ClientStream服务端Recv()到 EOF 后SendAndClose()的写法差异在示例工程的 handler 与文档中都有对应实现可按相同路径扩展。【免费下载链接】gofrAn opinionated GoLang framework for accelerated microservice development. Built in support for databases and observability.项目地址: https://gitcode.com/GitHub_Trending/go/gofr创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考