NestJS向Kafka发布消息调用emit方法报空引用错误排查
NestJS 调用Kafka客户端emit方法报null引用错误修复
问题表现
NestJS项目中Kafka消息消费逻辑运行正常,但需要将消费到的消息转发到第二个主题、或主动向Kafka发布消息时,抛出如下错误:
[Nest] 87277 - 07/08/2022, 2:45:34 PM ERROR [RpcExceptionsHandler] Cannot read properties of null (reading 'emit') TypeError: Cannot read properties of null (reading 'emit') at AppController.getHello (/path/src/app.controller.ts:31:17) at AppController.handleForeignData (/path/src/app.controller.ts:47:10) at /path/node_modules/@nestjs/microservices/context/rpc-context-creator.js:44:33 at processTicksAndRejections (node:internal/process/task_queues:96:5) at /path/node_modules/@nestjs/microservices/context/rpc-proxy.js:11:32 at ServerKafka.handleEvent (/path/node_modules/@nestjs/microservices/server/server.js:71:32) at Runner.processEachMessage (/path/node_modules/kafkajs/src/consumer/runner.js:231:9) at onBatch (/path/node_modules/kafkajs/src/consumer/runner.js:447:9) at Runner.handleBatch (/path/node_modules/kafkajs/src/consumer/runner.js:461:5) at /path/node_modules/kafkajs/src/consumer/worker.js:29:9
调用官方文档说明的emit()、send()两个发布方法都会触发相同报错,原业务代码如下:
import { Controller, Get, OnModuleInit } from '@nestjs/common'; import { AppService } from './app.service'; import { Client, ClientKafka, ClientProxy, EventPattern, MessagePattern, Payload } from '@nestjs/microservices' import { microserviceConfig } from "./microserviceConfig"; @Controller() export class AppController implements OnModuleInit { constructor(private readonly appService: AppService) {} @Client(microserviceConfig) client: ClientKafka; onModuleInit() { const requestPatterns = [ 'foreign-data', '-announcements' ]; requestPatterns.forEach(pattern => { this.client.subscribeToResponseOf(pattern); }); } @Get() getHello(): string { // fire event to kafka this.client.emit<string>('entity-created', 'some entity ' + new Date()); return this.appService.getHello(); } @EventPattern('foreign-data') async handleForeignData(payload: any) { // this.getHello(); console.log(JSON.stringify(payload) + ' created'); } }
根因分析
报错本质是执行emit调用时,this.client的值为null,由两个问题共同或单独导致:
- 未主动触发Kafka客户端连接:
@Client()装饰器声明的微服务客户端默认是懒加载模式,不会在应用启动阶段自动建立连接,未执行连接逻辑前客户端实例不会完成初始化,直接调用发送方法就会出现null引用。原代码仅在生命周期钩子中调用了subscribeToResponseOf,该方法仅用于请求-响应模式下的响应主题订阅,不会触发客户端连接。 - 类方法上下文丢失:
@EventPattern/@MessagePattern注册的消息处理器由NestJS RPC代理层调用,普通类方法被代理调用时如果没有做上下文绑定,this指向会脱离当前Controller类实例,此时访问this.client就会拿到null。
修复步骤
- 在
onModuleInit生命周期中显式调用客户端连接方法,等待连接建立完成后再执行业务逻辑。注意subscribeToResponseOf仅在需要使用send()方法做同步请求-响应交互时才需要配置,纯异步事件发布(用emit())不需要调用该方法。
修正后的生命周期代码:async onModuleInit() { // 主动建立Kafka生产者连接,等待实例初始化完成 await this.client.connect(); // 仅请求-响应模式需要保留以下订阅逻辑,纯发事件可删除 const requestPatterns = [ 'foreign-data', '-announcements' ]; requestPatterns.forEach(pattern => { this.client.subscribeToResponseOf(pattern); }); } - 避免消息处理器方法的上下文丢失:如果需要在消费逻辑中调用类内其他方法、或直接访问类属性,不要通过解耦的普通方法间接调用,优先在处理器方法中直接操作客户端实例;如果需要抽离公共方法,用箭头函数定义方法固定
this指向,或在构造函数中手动绑定方法上下文。
修正后的消息转发逻辑示例:@EventPattern('foreign-data') async handleForeignData(payload: any) { console.log(`${JSON.stringify(payload)} received`); // 直接在处理器中调用客户端发送,确保this指向当前Controller实例 this.client.emit<string>('entity-created', `forwarded payload: ${JSON.stringify(payload)} ${new Date()}`); } - 补充校验:确认
microserviceConfig配置中brokers地址列表、clientId参数填写正确,开启认证的Kafka集群需要同步配置SSL/SASL参数,避免连接建立失败。该类配置错误不会触发当前null引用报错,会直接抛出连接超时、认证失败类异常。
内容的提问来源于stack exchange,提问作者John O
相关产品推荐
相关产品推荐

