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

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模块:

  1. 安装依赖:
npm install @nestjs/terminus @nestjs/common
  1. 创建健康检查控制器:
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'],
      }),
    ]);
  }
}
  1. 启动代码和混合应用模式一致,确保HTTP服务和Kafka微服务同时运行即可。

结论

两种方案都能实现同时监听3000端口和订阅Kafka主题的需求:

  • 混合应用模式快速简单,适合基础的探针需求;
  • Terminus模块提供标准化的健康检查能力,适合生产环境中需要多维度状态校验的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 00:32:17