gRPC-Go中命令总线首次调用正常,后续请求失效的问题
Go gRPC + 六边形架构:首次请求正常后续无响应排查
问题概述
我正在开发首个Go语言项目,采用六边形架构与gRPC,遇到一个棘手问题:调用某个gRPC方法时首次执行正常,但后续请求均无法抵达处理器,连续5次请求仅首次生效。重启应用后,所有积压的命令会依次执行直至队列清空。仅在服务器Docker环境中出现此问题,本地Docker环境无异常。
核心代码片段
gRPC端口实现
// defining the command to send to the command bus protoMessage := productv1.CreateProductCommand{ Id: productID, ProductInfoId: in.Product.ProductInfoId, Product: in.Product, } // creating properly the message and validating the possible errors msg, err := messageManager.CreateMessage(ctx, &protoMessage) if err != nil { logger.WithContext(ctx).Error(fmt.Sprintf("CreateProduct: error create message %v", err)) return nil, err } // and finally, sending the message for the command bus err = s.commandBus.Send(ctx, finalMessage) if err != nil { logger.WithContext(ctx).Error(fmt.Sprintf("CreateProduct: error send message %v", err)) return nil, err }
命令处理器代码
func (h CreateProductHandler) Handle(ctx context.Context, c interface{}) error { logger.WithContext(ctx).Info(fmt.Sprintf("%s invoked with command = %v", h.HandlerName(), c)) eventId := uuid.NewString() errs := make([]*statuspb.Status, 0) if h.eventBus == nil || h.productRepo == nil || h.productInfoRepo == nil { err := errors.New("CreateProductHandler not properly initialized. Use NewCreateProductHandler") return err } cmd := c.(*productv1.CreateProductCommand) // after this cmd variable, comes the logic for product register, that is working as expected. }
注:处理器在错误时返回错误,正常时返回nil,怀疑返回nil是否存在问题。
处理器注册代码
productv1, err := product.NewApplicationFromSettings(settings) if err != nil { return nil, err } cqrsFacade, err := cqrs.NewFacade(newFacadeConfig( settings, redisstream.JSONMarshaler{}, *publisher, *subscriber, watermillLogger, router, func(commandBus *cqrs.CommandBus, eventBus *cqrs.EventBus) []cqrs.CommandHandler { commandHandlers := []cqrs.CommandHandler{} // previous command handlers... commandHandlers = append(commandHandlers, productv1.CommandHandlers(commandBus, eventBus)...) // other command handlers.... return commandHandlers }, ))
排查方向建议
- 命令总线阻塞检查:排查
commandBus.Send的实现,确认服务器环境下是否存在同步阻塞、资源未释放的情况,比如无缓冲通道未被消费、goroutine启动失败。 - 上下文生命周期验证:检查gRPC请求的
ctx是否在首次请求后被意外取消或超时,确认上下文在传递过程中未被截断,处理器中是否正确处理了上下文的取消信号。 - 资源限制排查:查看服务器Docker容器的资源配额(CPU、内存、文件句柄),检查是否存在数据库连接池耗尽、Redis连接未释放等情况,导致后续请求无法获取资源。
- 处理器并发模型验证:确认命令总线的并发配置,是否仅启动了单个消费goroutine,且首次请求的处理逻辑存在未退出的阻塞操作(比如死循环、未响应的外部RPC调用)。重点检查处理器中
// after this cmd variable后的注册逻辑,是否存在未完成的IO操作或阻塞。 - Watermill配置差异核对:对比服务器与本地环境的Watermill CQRS配置,特别是Redis Stream的消费者组、消息确认机制,确认是否存在未确认的消息导致消费阻塞。
- 返回值影响验证:处理器返回nil本身符合规范,但需确认是否在返回前完成了所有异步操作的清理,或命令总线是否需要特定返回值来确认消息处理完成。
内容的提问来源于stack exchange,提问作者Andre Cardoso
相关产品推荐
相关产品推荐

