寻求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
相关产品推荐
相关产品推荐

