NestJS Kafka微服务如何同时监听端口3000与Kafka主题?
NestJS Kafka微服务适配K8s存活/就绪探针方案
问题背景
你有一个基于Kafka传输的NestJS微服务,启动代码如下:
const app = await NestFactory.createMicroservice<MicroserviceOptions>(AppModule, { transport: Transport.KAFKA, options: { client: { brokers: ['localhost:9092'], } } }); await app.listen()
部署在Kubernetes上时,需要监听端口3000用于存活和就绪探针,但createMicroservice创建的实例无法直接设置监听端口,Kafka传输配置也不支持设置端口,想知道最佳解决方案,以及能否实现同时监听端口3000并订阅Kafka主题。
最佳解决方案:创建混合应用(同时支持HTTP服务+Kafka微服务)
NestJS原生支持混合应用模式,既能启动HTTP服务用于K8s探针检测,同时保留Kafka微服务的消息订阅能力,这是最直接的解决方案。
修改后的启动代码
import { NestFactory } from '@nestjs/core'; import { Transport } from '@nestjs/microservices'; import { AppModule } from './app.module'; async function bootstrap() { // 先创建HTTP应用,用于承载K8s探针接口 const app = await NestFactory.create(AppModule); // 连接Kafka微服务,配置消费者参数 app.connectMicroservice({ transport: Transport.KAFKA, options: { client: { brokers: ['localhost:9092'], }, consumer: { // 配置消费者组ID,必填项 groupId: 'your-service-consumer-group', // 订阅目标Kafka主题 subscribe: { topics: ['target-topic-name'], fromBeginning: true, }, }, }, }); // 启动所有微服务(包括Kafka消费者) await app.startAllMicroservices(); // 启动HTTP服务监听3000端口,供K8s探针调用 await app.listen(3000); } bootstrap();
添加健康检查接口
在AppController中新增一个简单的健康检查接口,供K8s探针访问:
import { Controller, Get } from '@nestjs/common'; @Controller() export class AppController { @Get('/health') healthCheck() { return { status: 'healthy', timestamp: new Date().toISOString() }; } }
之后K8s的存活/就绪探针就可以配置为访问http://<pod-ip>:3000/health,通过返回状态判断服务状态。
进阶方案:使用Terminus模块做专业健康检查
如果需要更规范、全面的健康检查(比如同时校验Kafka连接状态),可以使用NestJS官方的@nestjs/terminus模块:
- 安装依赖:
npm install @nestjs/terminus @nestjs/common
- 创建健康检查控制器:
import { Controller, Get } from '@nestjs/common'; import { HealthCheck, HealthCheckService, KafkaHealthIndicator } from '@nestjs/terminus'; @Controller('health') export class HealthController { constructor( private readonly healthCheckService: HealthCheckService, private readonly kafkaHealthIndicator: KafkaHealthIndicator, ) {} @Get() @HealthCheck() async check() { return this.healthCheckService.check([ // 校验Kafka连接状态 () => this.kafkaHealthIndicator.checkHealth('kafka', { brokers: ['localhost:9092'], }), ]); } }
- 启动代码和混合应用模式一致,确保HTTP服务和Kafka微服务同时运行即可。
结论
两种方案都能实现同时监听3000端口和订阅Kafka主题的需求:
- 混合应用模式快速简单,适合基础的探针需求;
- Terminus模块提供标准化的健康检查能力,适合生产环境中需要多维度状态校验的场景。
内容的提问来源于stack exchange,提问作者Yuval
相关产品推荐
相关产品推荐

