NestJS Kafka微服务全局异常过滤器无法捕获抛出异常
你遇到的全局异常过滤器不生效问题,核心是NestJS混合应用的默认配置隔离机制导致的:
你当前通过NestFactory.create()先创建HTTP应用实例,再调用connectMicroservice()挂载Kafka微服务的写法属于混合应用模式,这种模式下HTTP应用上注册的所有全局组件(过滤器、管道、守卫、拦截器)默认不会同步到挂载的微服务上下文。
你之前尝试的三种过滤器注册方式本身写法没有问题,但这些过滤器都只绑定到了HTTP服务层,根本不会作用于Kafka事件处理流程,自然捕获不到@EventPattern处理函数抛出的错误,最终直接触发Node的UnhandledPromiseRejection报错。
你可以根据自己的部署场景二选一:
方案1:保留现有混合应用结构,开启配置继承
如果确实需要同时保留HTTP服务和Kafka微服务能力,只需要在调用connectMicroservice()时传入第二个配置参数,开启inheritAppConfig开关,让微服务继承主应用的所有全局组件配置即可。
修改后的微服务连接代码如下:
app.connectMicroservice({ transport: Transport.KAFKA, options: { client: { clientId: SECRET_VALUE, brokers: [SECRET_HOST_ADDRESS], ssl: true, sasl: SOME_BOOLEAN_VALUE ? { mechanism: 'plain', username: SECRET_VALUE, password: SECRET_VALUE, } : undefined, }, consumer: { allowAutoTopicCreation: false, groupId: SECRET_VALUE, }, }, }, { inheritAppConfig: true }); // 新增这行配置即可
加完这个配置后,你之前注册的全局异常过滤器就会自动绑定到Kafka微服务上下文,可以正常捕获事件处理函数抛出的同步、异步异常。
方案2:改成纯微服务初始化(更适配Kafka Worker场景)
既然你的Kafka Worker节点不对外暴露任何REST接口,完全没必要启动冗余的HTTP服务监听端口,直接创建纯微服务实例即可。这种模式下不存在混合应用的配置隔离问题,全局组件的注册逻辑和普通HTTP应用完全一致,不会出现过滤器不生效的问题。
修改后的main.ts初始化逻辑参考:
async function bootstrap() { // 直接创建Kafka微服务实例,不再初始化HTTP应用 const app = await NestFactory.createMicroservice(KafkaWorkerAppModule, { transport: Transport.KAFKA, options: { client: { clientId: SECRET_VALUE, brokers: [SECRET_HOST_ADDRESS], ssl: true, sasl: SOME_BOOLEAN_VALUE ? { mechanism: 'plain', username: SECRET_VALUE, password: SECRET_VALUE, } : undefined, }, consumer: { allowAutoTopicCreation: false, groupId: SECRET_VALUE, }, }, logger: ['error', 'warn', 'debug', 'log', 'verbose'], }); // 全局组件直接在微服务实例上注册即可生效 app.useGlobalPipes( new ValidationPipe({ disableErrorMessages: false, whitelist: true, transform: true, }), ); app.useGlobalFilters(new KafkaWorkerExceptionFilter()); const logger: AppLogger = new AppLogger('Bootstrap'); const config: ConfigService = app.get(ConfigService); await app.listen(); logger.log(`Kafka Worker started`); logger.log(`Environment: ${config.nodeEnv}`); } bootstrap();
这种写法更轻量,不会占用多余端口,也完全符合纯事件消费节点的部署定位。
如果调整完配置后异常可以正常捕获,建议你在异常过滤器里补充偏移量提交的相关逻辑:如果Kafka消费者配置了手动提交偏移量,事件处理抛错时需要根据业务场景决定是重试、死信队列投递还是直接提交偏移量,避免消息重复消费或者丢失。
内容的提问来源于stack exchange,提问作者Maor Barazani

