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

复杂系统架构下与Async服务的通信方案技术问询

解决方案:与Async服务通信的Middleware及消息转换实现

核心思路

你的场景中,Middleware的核心作用是透明衔接Convertor与业务微服务:一边把Async发来的AmqpMessage转成业务能处理的格式,另一边把业务处理结果转回AmqpMessage发回Async。不用纠结Async服务本身,重点聚焦Middleware的拦截、转换、回传全链路逻辑。

Middleware实现步骤(以gRPC微服务为例)

1. 请求拦截与解析Middleware

在gRPC拦截器链中新增一层专门处理Amqp消息的Middleware,负责从队列获取Async的请求并转换:

  • 监听Async服务指定的请求队列
  • 收到AmqpMessage后,调用Convertor将其转成对应gRPC请求对象
  • 把转换后的请求传递给后续业务处理逻辑

示例伪代码:

// gRPC 请求拦截器Middleware
func AmqpRequestInterceptor(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
    // 从AMQP队列接收Async发来的请求消息
    amqpMsg, err := amqpClient.ReceiveMessage("async-request-queue")
    if err != nil {
        return nil, err
    }

    // 调用Convertor转换为gRPC请求
    grpcReq, err := convertor.AmqpToGrpc(amqpMsg)
    if err != nil {
        return nil, err
    }

    // 将请求传递给业务处理逻辑
    resp, err := handler(ctx, grpcReq)
    if err != nil {
        return nil, err
    }

    return resp, nil
}

2. 结果回传Middleware

业务处理完成后,拦截响应并转回AmqpMessage发回Async:

  • 拦截gRPC业务响应
  • 调用Convertor将响应转成AmqpMessage,必须携带原请求的CorrelationID(Async用来匹配原请求)
  • 将消息发送到Async监听的结果队列;若业务处理失败,需生成错误类型的AmqpMessage回传

示例伪代码:

// gRPC 响应回传Middleware
func AmqpResponseInterceptor(ctx context.Context, resp interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
    // 先执行业务处理逻辑
    result, err := handler(ctx, resp)
    if err != nil {
        // 转换错误信息为AmqpMessage并回传
        errorMsg := convertor.GrpcErrToAmqp(err)
        errorMsg.CorrelationID = ctx.Value("request-correlation-id").(string)
        amqpClient.SendMessage("async-response-queue", errorMsg)
        return nil, err
    }

    // 转换正常响应为AmqpMessage
    amqpResp := convertor.GrpcToAmqp(result)
    amqpResp.CorrelationID = ctx.Value("request-correlation-id").(string)
    amqpClient.SendMessage("async-response-queue", amqpResp)

    return result, nil
}

3. Convertor的解耦设计

将Convertor做成独立组件,Middleware仅依赖其接口,避免转换逻辑与链路逻辑耦合:

  • 定义统一转换接口AmqpConvertor,包含AmqpToGrpc、GrpcToAmqp、GrpcErrToAmqp三个核心方法
  • 各微服务根据自身业务请求/响应结构实现该接口,Middleware无需关心具体转换细节

关键注意事项

  • 关联ID追踪:全程传递Async消息自带的CorrelationID,确保Async能准确匹配请求与结果
  • 消息可靠性:Middleware需实现消息确认(ACK)、重试机制,避免消息丢失;异常场景必须回传错误消息,不能吞掉异常
  • 适配不同服务类型:如果是HTTP微服务,改用HTTP中间件(如Java Filter、Go net/http中间件)实现拦截与转换,逻辑与gRPC拦截器一致

内容的提问来源于stack exchange,提问作者Ali Mohammadnezhad

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 23:00:07