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

寻求NestJS Kafka微服务适配Avro Schema的实现示例

NestJS Kafka 结合 Avro Schema 实现示例

1. 安装依赖

首先安装KafkaJS的Confluent Schema Registry客户端:

npm install @kafkajs/confluent-schema-registry
# 或使用yarn
yarn add @kafkajs/confluent-schema-registry

2. 创建Avro Schema处理工具类

新建src/utils/avro-schema.util.ts,封装Schema Registry连接、序列化与反序列化逻辑:

import { SchemaRegistry } from '@kafkajs/confluent-schema-registry';
import { Injectable } from '@nestjs/common';

@Injectable()
export class AvroSchemaUtil {
  private registry: SchemaRegistry;

  constructor() {
    // 初始化Schema Registry客户端,按需调整配置
    this.registry = new SchemaRegistry({
      host: process.env.SCHEMA_REGISTRY_URL,
      ssl: true,
      // 若Schema Registry开启认证,添加以下配置
      // auth: {
      //   username: process.env.SCHEMA_REGISTRY_USER,
      //   password: process.env.SCHEMA_REGISTRY_PASSWORD,
      // },
    });
  }

  // 消费者用:反序列化Kafka二进制消息
  async deserializeMessage(message: any): Promise<any> {
    if (!message.value) return null;
    return this.registry.decode(message.value);
  }

  // 生产者用:序列化消息为Avro格式(可选)
  async serializeMessage(schemaId: number, payload: any): Promise<Buffer> {
    return this.registry.encode(schemaId, payload);
  }
}

3. 自定义Kafka序列化/反序列化器

新建src/utils/kafka-avro.serializer.ts,实现NestJS Kafka的序列化器接口:

import { Deserializer, Serializer } from '@nestjs/microservices';
import { AvroSchemaUtil } from './avro-schema.util';

// 消费者反序列化器
export class KafkaAvroDeserializer implements Deserializer {
  constructor(private readonly avroSchemaUtil: AvroSchemaUtil) {}

  async deserialize(value: any) {
    return this.avroSchemaUtil.deserializeMessage({ value });
  }
}

// 生产者序列化器(按需添加)
export class KafkaAvroSerializer implements Serializer {
  constructor(private readonly avroSchemaUtil: AvroSchemaUtil) {}

  async serialize(value: any) {
    // 需提前在Schema Registry注册Schema并获取ID,此处从环境变量读取
    const schemaId = parseInt(process.env.USER_EVENT_SCHEMA_ID);
    return this.avroSchemaUtil.serializeMessage(schemaId, value);
  }
}

4. 修改Kafka微服务配置

更新原有配置,替换默认的JSON序列化/反序列化器:

import { Transport } from '@nestjs/microservices';
import { AvroSchemaUtil } from './src/utils/avro-schema.util';
import { KafkaAvroDeserializer, KafkaAvroSerializer } from './src/utils/kafka-avro.serializer';

const avroSchemaUtil = new AvroSchemaUtil();

export const kafkaConfig = {
  transport: Transport.KAFKA,
  options: {
    client: {
      clientId: process.env.KAFKA_CLIENT_ID,
      ssl: true,
      brokers: process.env.KAFKA_BROKERS.split(','),
    },
    consumer: {
      groupId: process.env.KAFKA_CONSUMER_GROUP_ID,
    },
    subscribe: {
      fromBeginning: false,
    },
    // 配置自定义序列化/反序列化器
    deserializer: new KafkaAvroDeserializer(avroSchemaUtil),
    serializer: new KafkaAvroSerializer(avroSchemaUtil), // 生产者需要则保留
  },
};

5. 更新消费者代码

确保UserEvent类型与Avro Schema定义完全匹配,消费者逻辑无需额外修改,反序列化器会自动将二进制数据转为对应类型:

@Controller(USER)
@UseInterceptors(ClassSerializerInterceptor)
export class UserController {
  constructor(private readonly service: UserService) {}

  @EventPattern(process.env.USER_TOPIC, Transport.KAFKA)
  async processUserEvent(data: UserEvent) {
    // data已自动反序列化为符合Avro Schema的UserEvent实例
    await this.service.handleEvent(data);
  }
}

注意事项

  • 需提前在Schema Registry中注册对应Avro Schema,确保USER_EVENT_SCHEMA_ID与注册后的Schema ID一致
  • UserEvent接口/类的字段需严格匹配Avro Schema定义,避免类型不兼容问题
  • 若Schema Registry有ACL权限控制,需在AvroSchemaUtil中配置对应的认证参数

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 04:01:18