NestJS消费Kafka消息为Buffer无法正确反序列化问题
NestJS对接外部Kafka消费消息乱码解决方案
问题描述
在main.ts与controller.ts中完成Kafka微服务配置对接外部集群后,消费topic-tx主题的消息时,控制台输出为Buffer格式乱码,无法反序列化为符合IResponseValue接口定义的结构化数据。
现有配置
main.ts 微服务定义
app.connectMicroservice<MicroserviceOptions>({ transport: Transport.KAFKA, options: { client: { brokers: ['localhost:9092'], }, consumer: { groupId: 'consumer-1', }, }, }); await app.startAllMicroservices(); app.useGlobalPipes( new ValidationPipe({ transform: true, }), );
controller.ts 消费逻辑
export class Controller implements OnModuleInit { constructor(private readonly service: Service) {} @Client({ transport: Transport.KAFKA, options: { client: { brokers: ['localhost:9092'], }, consumer: { groupId: 'consumer-1', }, }, }) client: ClientKafka; async onModuleInit() { this.client.subscribeToResponseOf('topic-tx'); await this.client.connect(); Logger.log( 'consumer assignments: ' + JSON.stringify(this.client.getConsumerAssignments()), ); } @MessagePattern('topic-tx', Transport.KAFKA) async handleEntityCreated(@Payload('value') message: IResponseValue) { console.log('Received event: ', message); } }
异常输出
Received event: $da4fa4c3-9e91-43d9-acaa-a07d1fed2635"� *0x5bcd9e419a11AB71f5eea1a2CFf9B0694C990Baf29*0xDFB50936C5d83b8367BDC01B17c386203AA60368*4611920x0:�0xd3fc98640000000000000000000000005bcd9e419a11ab71f5eea1a2cff9b0694c990baf00000000000000000000000000000000000000000000000000000000000000380000000000000000000000000000000000000000000000000000000000000060000000000000000000000000000000000000000000000000000000000000004034346532343639353263666337363430663535306162386233343363633564633865623033303937373534343030343935306134646432393134663864663561B�0xf901261d8082b42794dfb50936c5d83b8367bdc01b17c386203aa6036880b8c4d3fc98640000000000000000000000005bcd9e419a11ab71f5eea1a2cff9b0694c990baf00000000000000000000000000000000000000000000000000000000000000380000000000000000000000000000000000000000000000000000000000000060000000000000000000000000000000000000000000000000000000000000004034346532343639353263666337363430663535306162386233343363633564633865623033303937373534343030343935306134646432393134663864663561820a96a0d524dbe02b8120452b988b0409281e0cf4f7db1ad1dba3c3acfdb11d18c8cc5fa043d1cdea80c15f180854e4de70051cff66f59067ef6a6bbb90156a30285d100eJB0x02dff5ee07b16609c959660789dd743d648a5f44f4ad3651fefe76f7ee004134�legacy*� B0x02dff5ee07b16609c959660789dd743d648a5f44f4ad3651fefe76f7ee004134B0x362711a341c6e4647f8978a9bb01df8097d9dd9177989389fc17627c66ee1615�%@R�0x00000000000000000000000000000000000000000000008000000000000000000000000000000000000020000000000000000040000000000000000000012000000000000000110000000020000000000000000000000000000000000000000000000000020000000000000000000800000000002000000000000000000000000000000000000000000000000000001000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000008000000000000000020000004000000000000000000000000000000000000000000000200080000000800Z� *0xDFB50936C5d83b8367BDC01B17c386203AA60368B0xc3d58168c5ae7397731d063d5bbf3d657854427343f4c083240f7aacaa2d0f62B0x0000000000000000000000005bcd9e419a11ab71f5eea1a2cff9b0694c990bafB0x0000000000000000000000000000000000000000000000000000000000000000B0x0000000000000000000000005bcd9e419a11ab71f5eea1a2cff9b0694c990baf�0x00000000000000000000000000000000000000000000000000000000000000380000000000000000000000000000000000000000000000000000000000000001"7TransferSingle(address,address,address,uint256,uint256)*2 from*0x0000000000000000000000000000000000000000*0 to*0x5bcd9e419a11AB71f5eea1a2CFf9B0694C990Baf* id56* value1*6 operator*0x5bcd9e419a11AB71f5eea1a2CFf9B0694C990Baf0�%:B0x02dff5ee07b16609c959660789dd743d648a5f44f4ad3651fefe76f7ee004134JB0x362711a341c6e4647f8978a9bb01df8097d9dd9177989389fc17627c66ee1615Z� *0xDFB50936C5d83b8367BDC01B17c386203AA60368B0x6bb7ff708619ba0610cba295a58592e0451dee2622938c8755667688daf3529bB0x0000000000000000000000000000000000000000000000000000000000000038�0x0000000000000000000000000000000000000000000000000000000000000020000000000000000000000000000000000000000000000000000000000000004034346532343639353263666337363430663535306162386233343363633564633865623033303937373534343030343935306134646432393134663864663561"URI(string,uint256)*I value@44e246952cfc7640f550ab8b343cc5dc8eb030977544004950a4dd2914f8df5a* id560�%:B0x02dff5ee07b16609c959660789dd743d648a5f44f4ad3651fefe76f7ee004134JB0x362711a341c6e4647f8978a9bb01df8097d9dd9177989389fc17627c66ee1615P`��h��r0x0� MetalToken�latest:devB$0afcfd48-c20b-4db7-ab35-73f90e937c37
根因分析
NestJS Kafka模块默认内置的KafkaDeserializer仅适配NestJS微服务自有通信协议:该协议会在生产端发送消息时包裹Nest自定义二进制头、默认做JSON序列化,消费端按固定规则解析协议头后再反序列化内容。
对接外部Kafka集群时,生产端不会遵循Nest私有协议封装消息,默认解析器会把原始二进制Buffer错误按Nest协议格式拆解,最终输出乱码。
另外当前配置存在重复消费者问题:main.ts中通过connectMicroservice已经启动了一个groupId为consumer-1的消费者,Controller中@Client装饰器又声明了同groupId的第二个消费者,会导致分区分配异常、重复消费问题。
解决步骤
1. 实现自定义反序列化器
根据生产端的消息编码格式,实现KafkaDeserializer接口完成消息解码:
- 如果生产端发送的是UTF-8编码的JSON消息,直接将Buffer转字符串后JSON解析即可
- 如果是Avro/Protobuf/自定义二进制格式(从日志看包含ERC1155链上事件签名,属于自定义二进制编码),需要引入对应格式的解码库(如ethers.js ABI解码器、avsc等)处理
以JSON格式消息为例,自定义反序列化器代码如下:
import { KafkaDeserializer } from '@nestjs/microservices'; import { ConsumerRecord } from '@nestjs/microservices/external/kafka.interface'; export class CustomKafkaDeserializer implements KafkaDeserializer { deserialize(message: ConsumerRecord): any { // 二进制Buffer转UTF-8字符串 const rawValue = message.value?.toString('utf8') ?? null; if (!rawValue) return { ...message, value: null }; try { // JSON解析为结构化对象 return { ...message, value: JSON.parse(rawValue), }; } catch (err) { // 解析失败时返回原始字符串和Buffer,方便排查格式问题 return { ...message, value: rawValue, rawBuffer: message.value, parseError: err.message, }; }
相关产品推荐
相关产品推荐

