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

NestJS网关与Spring Boot微服务NATS请求响应通信异常排查

NestJS与Spring Boot通过NATS请求响应模式通信超时排查

问题背景

我正在实现一套以NATS作为消息中间件的微服务系统,采用NestJS网关与Spring Boot微服务进行请求-响应模式通信。当前出现的问题是:NestJS网关发送请求后,Spring Boot微服务能正常接收、处理并返回响应,但NestJS始终无法收到该响应,最终触发超时。日志显示Spring Boot已向正确的replyTo地址发布响应,但NestJS端无任何接收迹象。

NestJS网关代码

import { Body, Controller, Inject, Post } from "@nestjs/common";
import { ClientProxy, RpcException } from "@nestjs/microservices";

import { SERVICES } from "src/config";
import { LoginUserRequestDto } from "./dtos/LoginUser.dto";
import { catchError, lastValueFrom, timeout } from "rxjs";

export class LoginUserRequestDto {
    username: string;
    password: string;
}

export enum SERVICES {
    NATS_SERVICE = "NATS_SERVICE",
}


@Controller("auth")
export class AuthController {
    constructor(
        @Inject(SERVICES.NATS_SERVICE) private readonly client: ClientProxy,
    ) {}

    @Post("/login")
    async login(@Body() loginUserRequestDto: LoginUserRequestDto) {
        console.log(loginUserRequestDto);

        const msResponse = await lastValueFrom(
            this.client.send("login.crud.ms", loginUserRequestDto).pipe(
                timeout(10000),
                catchError((error) => {
                    throw new RpcException(error);
                }),
            ),
        );

        console.log(msResponse);
        return msResponse;

        // return this.client.send("login.crud.ms", loginUserRequestDto).pipe(
        //  timeout(5000),
        //  catchError((error) => {
        //      throw new RpcException(error);
        //  }),
        // );
    }
}

NestJS NATS客户端配置

import { Module } from "@nestjs/common";
import { ClientsModule, Transport } from "@nestjs/microservices";

import { enviromentVariables, SERVICES } from "src/config";

@Module({
    imports: [
        ClientsModule.register([
            {
                name: SERVICES.NATS_SERVICE,
                transport: Transport.NATS,
                options: {
                    servers: enviromentVariables.natsSever,
                },
            },
        ]),
    ],
    exports: [
        ClientsModule.register([
            {
                name: SERVICES.NATS_SERVICE,
                transport: Transport.NATS,
                options: {
                    servers: enviromentVariables.natsSever,
                },
            },
        ]),
    ],
})
export class TransportsModule {}

Spring Boot微服务代码

@RestController
@Validated
@RequestMapping("/users")
public class UserController {

    private UserService userService;
    private final Connection natsConnection;
    private final ObjectMapper objectMapper;

    UserController(UserService userService, Connection natsConnection, ObjectMapper objectMapper) {
        this.userService = userService;
        this.natsConnection = natsConnection;
        this.objectMapper = objectMapper;
    }

    @PostConstruct
    public void subscriberTopic() {
        loginUserMS();
    }

    public void loginUserMS() {
        natsConnection.createDispatcher(msg -> {
            try {
                String rawJSON = new String(msg.getData(), StandardCharsets.UTF_8);
                System.out.println(rawJSON); // 控制台输出: {"pattern":"login.crud.ms","data":{"username":"pedro123","password":"1234"},"id":"23f2bd6f679cc394c433d"}

                
                JsonNode rootNode = objectMapper.readTree(rawJSON);

                JsonNode dataNode = rootNode.get("data");
                String correlationId = rootNode.get("id").asText();

                
                LoginUserRequestDto loginUserRequestDto = objectMapper.treeToValue(dataNode,
                        LoginUserRequestDto.class);

               

                // 调用业务服务
                UserDto loggedUser = userService.login(loginUserRequestDto);


                System.out.println(loggedUser.toString());

                Map<String, Object> replyWrapper = new HashMap<>();
                replyWrapper.put("id", correlationId);
                replyWrapper.put("err", null);
                replyWrapper.put("response", loggedUser);

                // 序列化响应
                String response = objectMapper.writeValueAsString(loggedUser);
                // String response = objectMapper.writeValueAsString(replyWrapper);


                System.out.println(response); // 控制台输出: {"id":"98e54e58-af64-417d-97ba-637e6be6e4b3","name":"Enrique","username":"pedro123","email":"pepe123@gmail.com"}

                // natsConnection.publish(msg.getReplyTo(),
                // response.getBytes(StandardCharsets.UTF_8));
                natsConnection.publish(msg.getReplyTo(), response.getBytes(StandardCharsets.UTF_8));
            } catch (Exception e) {
                e.printStackTrace();
                if (msg.getReplyTo() != null) {
                    String errorResponse = "{\"error\": \"" + e.getMessage() + "\"}";
                    natsConnection.publish(msg.getReplyTo(), errorResponse.getBytes(StandardCharsets.UTF_8));
                }
                throw new Error(e.getMessage());
            }
        }

        ).subscribe("login.crud.ms");
    }

问题排查与解决方案

核心原因

NestJS的NATS客户端在请求-响应模式下,期望收到的响应必须包含与请求对应的id字段,用于关联请求和响应。当前Spring Boot返回的是原始的UserDto对象,没有包裹NestJS要求的响应格式,导致NestJS无法识别这是对应请求的响应,因此一直等待直到超时。

修复步骤

  1. 修改Spring Boot的响应格式:使用代码中已定义的replyWrapper对象进行序列化,而不是直接返回UserDto。取消注释以下代码:

    // String response = objectMapper.writeValueAsString(loggedUser);
    String response = objectMapper.writeValueAsString(replyWrapper);
    

    这样返回的响应结构会包含id(关联请求的correlationId)、err(错误信息,成功时为null)、response(实际业务数据),完全匹配NestJS的期望格式。

  2. 统一错误响应格式:异常处理中返回的响应也需要调整为符合NestJS要求的结构,避免出现格式不匹配导致的超时:

    Map<String, Object> errorWrapper = new HashMap<>();
    errorWrapper.put("id", correlationId);
    errorWrapper.put("err", e.getMessage());
    errorWrapper.put("response", null);
    String errorResponse = objectMapper.writeValueAsString(errorWrapper);
    
  3. 验证网络与NATS配置:确认NestJS和Spring Boot连接的NATS服务器地址完全一致,且没有防火墙或网络策略阻止replyTo主题的消息传输。可以使用NATS命令行工具(如nats sub <replyTo主题>)订阅对应主题,验证Spring Boot是否确实发送了响应。

验证修复

修改后,NestJS将能通过响应中的id字段匹配到对应的请求,结束等待并返回业务数据,解决超时问题。


内容的提问来源于stack exchange,提问作者Héctor Romero

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 10:47:09