NestJS微服务Gateway与Users通过RabbitMQ通信时连接关闭报错
NestJS微服务通过RabbitMQ通信时出现"Connection closed"错误
问题背景
我拥有gateway和users两个微服务,预期通过RabbitMQ实现相互通信。
gateway服务代码
// app.module.ts import { Module } from '@nestjs/common'; import { AppController } from './app.controller'; import { AppService } from './app.service'; import { ClientsModule } from '@nestjs/microservices'; import { UsersController } from './users/users.controller'; @Module({ imports: [ ClientsModule.register([ { name: 'USERS_SERVICE', options: { urls: ['amqp://localhost:5672/hello'], queue: 'users_queue', queueOptions: { durable: false }, }, }, ]), ], controllers: [AppController, UsersController], providers: [AppService], }) export class AppModule {}
// users.controller.ts import { Controller, Post, Body, Get, Param, Patch, Delete, Inject } from '@nestjs/common'; import { ClientProxy } from '@nestjs/microservices'; import { CreateUserDto, UpdateUserDto } from './dto'; @Controller('users') export class UsersController { constructor(@Inject('USERS_SERVICE') private usersClient: ClientProxy) {} @Post() create(@Body() createUserDto: CreateUserDto): Observable<object> { return this.usersClient.send('createUser', createUserDto); } @Get() findAll(): Observable<object> { return this.usersClient.send('findAllUsers', {}); } @Get(':id') findOne(@Param('id') id: string) { return this.usersClient.send('findOneUser', id); } @Patch(':id') update(@Param('id') id: string, @Body() updateUserDto: UpdateUserDto) { return this.usersClient.send('updateUser', { id, updateUserDto }); } @Delete(':id') remove(@Param('id') id: string) { return this.usersClient.send('removeUser', id); } }
users服务代码
// main.ts import { NestFactory } from '@nestjs/core'; import { AppModule } from './app.module'; import { MicroserviceOptions, Transport } from '@nestjs/microservices'; async function bootstrap() { const app = await NestFactory.createMicroservice<MicroserviceOptions>( AppModule, { transport: Transport.RMQ, options: { urls: ['amqp://localhost:5672/hello'], queue: 'users_queue', queueOptions: { durable: false }, }, }, ); app.listen(); } bootstrap();
// users.controller.ts import { Controller } from '@nestjs/common'; import { MessagePattern } from '@nestjs/microservices'; import { UsersService } from './users.service'; @Controller() export class UsersController { constructor(private readonly usersService: UsersService) {} @MessagePattern('findAllUsers') findAll() { return this.usersService.findAll(); } }
// users.service.ts import { Injectable } from '@nestjs/common'; @Injectable() export class UsersService { findAll() { return `This action returns all users`; } }
错误现象
发送GET请求http://localhost:3000/users时,gateway服务抛出以下错误:
[Nest] 21343 - 11/01/2023, 10:24:17 PM ERROR [ExceptionsHandler] Connection closed Error: Connection closed at ClientTCP.handleClose (/home/arta/project/vsCode/nestjs/microHub/gateway/node_modules/.pnpm/@nestjs+microservices@10.2.7_@nestjs+common@10.0.0_@nestjs+core@10.0.0_amqp-connection-manage_a4katzriqmc2wl6sp3kq4c2z64/node_modules/@nestjs/microservices/client/client-tcp.js:102:25) at Socket. (/home/arta/project/vsCode/nestjs/microHub/gateway/node_modules/.pnpm/@nestjs+microservices@10.2.7_@nestjs+common@10.0.0_@nestjs+core@10.0.0_amqp-connection-manage_a4katzriqmc2wl6sp3kq4c2z64/node_modules/@nestjs/microservices/client/client-tcp.js:92:55) at Socket.emit (node:events:526:35) at TCP. (node:net:337:12) ^C ELIFECYCLE Command failed.
注:RabbitMQ已正常运行,预期返回内容为"This action returns all users"。
问题原因
错误栈中明确提到ClientTCP,说明gateway服务的ClientsModule默认使用了TCP传输协议,而非RabbitMQ对应的RMQ传输协议。在gateway的ClientsModule.register配置中,没有指定transport: Transport.RMQ,导致NestJS尝试用TCP去连接RabbitMQ服务,自然无法建立有效连接,最终触发"Connection closed"错误。
解决方案
在gateway服务的ClientsModule.register配置中添加transport: Transport.RMQ,确保和users服务的传输方式一致。修改后的app.module.ts代码如下:
// app.module.ts import { Module } from '@nestjs/common'; import { AppController } from './app.controller'; import { AppService } from './app.service'; import { ClientsModule, Transport } from '@nestjs/microservices'; import { UsersController } from './users/users.controller'; @Module({ imports: [ ClientsModule.register([ { name: 'USERS_SERVICE', transport: Transport.RMQ, // 新增此行,指定RMQ传输 options: { urls: ['amqp://localhost:5672/hello'], queue: 'users_queue', queueOptions: { durable: false }, }, }, ]), ], controllers: [AppController, UsersController], providers: [AppService], }) export class AppModule {}
修改后重启gateway服务,再次发送请求即可正常获取预期结果。
内容的提问来源于stack exchange,提问作者arta
相关产品推荐
相关产品推荐

