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

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,
      };
    }
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 13:39:18