- 后端
- 微服务
- 云原生
【免费下载链接】midway
🍔 A Node.js Serverless Framework for front-end/full-stack developers. Build the application for next decade. Works on AWS, Alibaba Cloud, Tencent Cloud and traditional VM/Container. Super easy integrate with React and Vue. 🌈
导读
本文以 midway 仓库中 @midwayjs/kafka 组件 的版本演进为脉络,结合其源码实现,系统讲解 Midway 框架中 Kafka 消息组件的完整用法:从消费者(Consumer)、生产者(Producer)到管理端(Admin)的配置与编程模型,以及该组件从引入、迭代到完善的演进路径。读者读完后将掌握如何基于@midwayjs/kafka在 Midway 应用中快速落地 Kafka 消息收发、如何配置连接参数、如何复用同一 Kafka 实例,并理解其底层基于 KafkaJS 的封装原理。
一、组件演进:从引入到稳定
@midwayjs/kafka是 midway 中面向 Kafka 场景的官方子包,其底层基于 kafkaJS 与 package.json 中"kafkajs": "2.2.4"依赖)。从 CHANGELOG 可以清晰还原该组件的发展轨迹:
| 版本 | 时间 | 类型 | 核心变更 |
|---|---|---|---|
| 3.4.0-beta.4 | 2022-07-04 | Feature | 新增 kafka 组件(#2062),同时修复 config export default 大小写问题(#2089) |
| 3.4.0-beta.12 | 2022-07-20 | Bug Fix | passport 兼容性代码调整(#2133) |
| 3.4.10 | 2022-08-12 | Bug Fix | 捕获 Kafka 启动错误(#2230),避免启动异常未被感知 |
| 3.4.11 | 2022-08-16 | Feature | 更新 kafka framework 并补充测试示例(#2236) |
| 3.6.0 | 2022-10-10 | Feature | 支持 guard(#2345),将守卫机制引入 kafka 消费链路 |
| 3.7.0 | 2022-10-29 | - | 随主版本发布,无额外变更 |
从版本记录可以看到:该组件于 3.4.0 时代作为全新特性被引入 midway,随后经历"启动错误捕获"的可靠性修复、"framework 更新与测试补充"的能力完善,再到 3.6.0 引入 guard 守卫能力,最终随主版本稳定发布。这些变更恰好对应了组件源码中「消费者资源初始化 → 运行 → 销毁」完整生命周期管理的可靠性设计,以及基于applyMiddleware的中间件/守卫执行链路(详见 framework.ts)。
二、组件架构:四层职责划分
@midwayjs/kafka的源码结构非常清晰,共 7 个源文件,各司其职:
packages/kafka/src/ ├── index.ts # 模块统一出口,重新导出全部能力 ├── configuration.ts # 组件配置入口,注册 namespace 与默认日志 ├── decorator.ts # @KafkaConsumer 装饰器定义 ├── framework.ts # MidwayKafkaFramework:消费者生命周期管理核心 ├── manager.ts # KafkaManager:Kafka 客户端实例注册表(单例) ├── service.ts # KafkaProducerFactory / KafkaAdminFactory:生产者与管理端工厂 └── interface.ts # 全部 TypeScript 类型定义2.1 配置入口(configuration.ts)
configuration.ts 通过@Configuration注册了kafka命名空间,并注入了默认配置:kafka: {}空对象作为初始配置,同时注册了一个名为kafkaLogger的日志客户端(fileLogName: 'midway-kafka.log'),用于独立记录 Kafka 相关日志。组件在onReady阶段主动实例化KafkaProducerFactory,确保应用就绪时生产者工厂已完成初始化。
2.2 装饰器与消费者声明(decorator.ts)
decorator.ts 定义了核心装饰器@KafkaConsumer(consumerName)。其内部通过saveModule注册模块、saveClassMetadata保存消费者名称,并自动赋予Scope(Request)请求作用域与@Provide()依赖注入能力。这意味着:
- 每个消费者类默认处于请求级作用域(每条消息到达时创建独立实例,天然隔离状态);
- 消费者名称(如
sub1)与配置项中的键名一一对应,用于绑定订阅配置。
2.3 框架核心(framework.ts)
framework.ts 中的MidwayKafkaFramework是整个消费链路的引擎,其run()方法完成以下关键流程:
- 扫描消费者:通过
DecoratorManager.listModule(KAFKA_DECORATOR_KEY)收集所有@KafkaConsumer装饰的类,建立「名称 → 类」映射表; - 创建资源:对每个消费者配置,若指定了
kafkaInstanceRef则复用已注册的 Kafka 实例(找不到时抛出MidwayCommonError),否则基于connectionOptions新建 Kafka 客户端并注册进KafkaManager;随后创建consumer、connect()、subscribe(); - 绑定回调:根据消费者类实现的是
eachBatch还是eachMessage方法自动选择运行模式,并包装为带链路追踪(tracing)与中间件的执行函数; - 启动与销毁:通过
resourceStart调用consumer.run(runConfig),在beforeStop阶段统一disconnect()。
值得关注的是链路追踪集成:每个消费回调都会通过MidwayTraceService.runWithEntrySpan包裹,携带midway.protocol: 'kafka'、midway.kafka.topic等属性,并支持通过kafka.tracing.extractor自定义从消息头提取链路上下文(详见 framework.ts)。
2.4 实例注册表(manager.ts)
manager.ts 中的KafkaManager是一个线程内单例,维护Map<string, Kafka>。它承担「共享实例」的关键职责:当多个消费者、生产者或管理端希望复用同一个 Kafka 连接时,通过kafkaInstanceRef引用同一名称即可,避免重复建连。
三、消费者:消息监听实战
3.1 声明式消费者
在业务代码中,通过@KafkaConsumer('消费者名')装饰类并实现eachMessage(逐条消息)或eachBatch(批量消息)方法即可:
import { Provide, Inject } from '@midwayjs/core'; import { KafkaConsumer, Context, IKafkaConsumer } from '@midwayjs/kafka'; import { EachMessagePayload } from 'kafkajs'; @Provide() @KafkaConsumer('sub1') export class UserConsumer implements IKafkaConsumer { @Inject() ctx: Context; async eachMessage(payload: EachMessagePayload) { const { topic, partition, message } = payload; this.ctx.logger.info( `topic: ${topic}, partition: ${partition}, value: ${message.value.toString()}` ); } }与仓库测试中的用法一致(见 index.test.ts),消费者的名称(sub1)必须与配置文件中的键名对应:
// config 或 globalConfig 中 kafka: { consumer: { sub1: { connectionOptions: { clientId: 'my-app', brokers: [process.env.KAFKA_URL || 'localhost:9092'], }, consumerOptions: { groupId: 'groupId-test-' + Math.random(), }, subscribeOptions: { topics: ['topic-test-1'], fromBeginning: false, }, }, }, }3.2 配置项详解
消费者相关的四组配置均定义于 interface.ts 的IKafkaConsumerInitOptions:
| 配置字段 | 类型 | 说明 |
|---|---|---|
connectionOptions | KafkaConfig | KafkaJS 的客户端连接参数,如clientId、brokers等,透传给new Kafka() |
consumerOptions | ConsumerConfig | KafkaJS 消费者参数,最常用的是groupId(消费组) |
subscribeOptions | ConsumerSubscribeTopics / ConsumerSubscribeTopic | 订阅参数,topics指定订阅主题列表,fromBeginning控制是否从最早 offset 消费 |
consumerRunConfig | ConsumerRunConfig | 运行参数,可覆盖默认的eachMessage/eachBatch行为 |
kafkaInstanceRef | string | 可选,指定复用已注册的 Kafka 实例名称(共享连接) |
从 framework.ts 的resourceInitialize实现可见:connectionOptions与logCreator(将 KafkaJS 日志级别映射到 midway 的kafkaLogger)一起构造 Kafka 实例;映射关系为NOTHING→none、ERROR→error、WARN→warn、INFO→info、DEBUG→debug(见 framework.ts)。
3.3 多消费者与共享实例
仓库测试提供了两种典型场景(index.test.ts):
- 多主题多消费者:分别声明
sub1、sub2两个消费者,各自订阅topic-test-1、topic-test-2,互不干扰; - 共享 Kafka 实例:
sub2配置kafkaInstanceRef: 'sub1'复用sub1的客户端连接,仅新建 consumer 与订阅关系,节省连接资源。
四、生产者:消息发送实战
生产者通过KafkaProducerFactory工厂管理(service.ts)。配置文件结构如下:
kafka: { producer: { clients: { producer1: { connectionOptions: { clientId: 'my-app', brokers: [process.env.KAFKA_BROKERS || 'localhost:9092'], }, producerOptions: { createPartitioner: Partitioners.DefaultPartitioner, // 可选分区策略 }, }, }, }, },在业务代码中通过依赖注入获取工厂并取得命名生产者:
import { Inject } from '@midwayjs/core'; import { KafkaProducerFactory } from '@midwayjs/kafka'; @Provide() export class OrderService { @Inject() producerFactory: KafkaProducerFactory; async sendMessage() { const producer = this.producerFactory.get('producer1'); await producer.send({ topic: 'order-topic', messages: [{ key: 'order-1', value: JSON.stringify({ id: 1 }) }], }); } }仓库测试展示了完整收发闭环(index.test.ts):通过app.getApplicationContext().getAsync(Kafka.KafkaProducerFactory)获取工厂,get('producer1')取得实例后send消息,再由原生 KafkaJS consumer 订阅验证消息到达。
生产者工厂的底层行为(service.ts)与消费者框架一致:支持kafkaInstanceRef复用实例;创建成功后监听producer.connect事件并记录日志;销毁时producer.disconnect()。此外,KafkaProducerFactory通过bindTraceContext对send/sendBatch做了包装(service.ts),自动为每条消息注入链路上下文到 headers 中,保证"生产-消费"跨进程链路可追踪。
五、管理端(Admin):主题与消费组管理
@midwayjs/kafka还提供了KafkaAdminFactory封装 KafkaJS 的 Admin 能力(service.ts):
kafka: { admin: { clients: { admin1: { connectionOptions: { clientId: 'my-app', brokers: [process.env.KAFKA_BROKERS || 'localhost:9092'], }, }, }, }, },使用方式与生产者类似:app.getApplicationContext().getAsync(Kafka.KafkaAdminFactory)后get('admin1'),即可调用createTopics、listTopics、deleteTopics、listGroups等管理操作。仓库测试覆盖了完整的"创建主题 → 校验存在 → 删除主题 → 校验消费组"流程(index.test.ts)。
六、本地开发与测试
组件仓库内置了 Kafka 本地启动脚本(scripts):
start.sh:一键拉起本地 Kafka 环境;stop.sh:停止本地 Kafka;kafka-group.yml:编排文件(供脚本调用)。
配合组件测试(index.test.ts),无需真实 broker 也能覆盖大部分逻辑;需要真实验证时,测试通过process.env.KAFKA_URL/process.env.KAFKA_BROKERS读取连接地址,默认回退到localhost:9092,并使用Math.random()生成随机groupId避免消费组冲突(fixtures/base-app/src/configuration.ts)。
测试用例目录中还包含了 entry-trace.test.ts 与 trace.test.ts,专门验证消费者入口 span 与生产者注入的链路上下文是否贯通,印证了上文所述的 tracing 集成能力。
七、常见问题与注意事项
- 消费者类必须实现
eachMessage或eachBatch:框架通过ClzProvider.prototype['eachBatch']是否存在来自动判定运行模式(framework.ts),两者都不实现将导致消费回调为空。 kafkaInstanceRef引用不存在会直接抛错:消费者、生产者、管理端三处均校验实例是否存在,错误信息形如kafka instance xxx not found(framework.ts),配置共享实例前需确认引用名称正确。- 生产与消费共用连接:推荐消费者先注册实例名,生产者 / 管理端通过
kafkaInstanceRef复用,既省连接又保证链路追踪上下文一致(参见 index.test.ts 的共享实例测试)。 - 日志独立记录:Kafka 组件日志输出到独立的
midway-kafka.log(configuration.ts),排查问题时可单独关注该文件。
结语
@midwayjs/kafka是 midway 官方对 Kafka 场景的完整封装:以@KafkaConsumer装饰器 + 配置驱动的方式隐藏了 KafkaJS 复杂的建连、订阅、运行细节,同时保留了对connectionOptions、consumerOptions、subscribeOptions等底层参数的透传能力;生产者与管理端则通过统一工厂模式提供即取即用的服务。结合组件演进历史中"启动错误捕获""guard 支持""测试补充"等迭代,可以确认这是一个持续打磨、面向生产环境的成熟组件。读者可结合本文配置示例与 组件源码 中的类型定义,快速接入自己的 midway 应用。
- 后端
- 微服务
- 云原生
【免费下载链接】midway
🍔 A Node.js Serverless Framework for front-end/full-stack developers. Build the application for next decade. Works on AWS, Alibaba Cloud, Tencent Cloud and traditional VM/Container. Super easy integrate with React and Vue. 🌈
相关推荐
.NET Aspire 集成 Apache Kafka:基于 Confluent.Kafka 的生产者与消费者组件实战指南
.NET Aspire 集成 Apache Kafka:基于 Confluent.Kafka 的生产者与消费者组件实战指南 Aspire.Confluent.K
云原生后端微服务可观测性开发工具Obsidian终极美化指南:3分钟打造你的个性化知识管理神器
Obsidian终极美化指南:3分钟打造你的个性化知识管理神器 你是否正在使用Obsidian进行知识管理,但总觉得界面不够个性化?想要让笔记应用既美观又高效吗
文档知识管理Midway v4 集成 Apollo GraphQL:`@midwayjs/apollo` 与 `@midwayjs/graphql` 双包架构实战指南
Midway v4 集成 Apollo GraphQL: @midwayjs/apollo 与 @midwayjs/graphql 双包架构实战指南 Midwa
后端微服务云原生
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考