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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 23:25:32