EvaQueue.js 是一个统一的高性能消息队列 API 抽象层,基于适配器模式设计。通过一套一致的接口操作 Kafka、阿里云 MNS 和 NATS JetStream,切换后端无需改动业务代码。
- 统一 API — Queue(队列)和 Topic(主题)两种模式,接口一致
- 按需安装 — MQ 客户端为 peer 依赖,只用你需要的后端
- 动态加载 — 适配器在首次使用时才
import(),无冗余初始化 - TypeScript 优先 — 完整类型定义,IDE 友好
- 优雅退出 — 内置
enableGracefulExit(),安全关闭消费者连接 - Pure ESM — 原生 ES Module,适配 Node.js 现代生态
npm install evamq按需安装消息队列客户端(至少一个):
# Kafka
npm install @confluentinc/kafka-javascript
# 阿里云 MNS
npm install ali-mns
# NATS JetStream
npm install @nats-io/transport-node @nats-io/jetstream注意:安装
@confluentinc/kafka-javascript如遇ld: symbol(s) not found for architecture x86_64,尝试:CPPFLAGS=-I/usr/local/opt/openssl/include LDFLAGS=-L/usr/local/opt/openssl/lib npm install
生产消息:
importMQfrom'evamq';importMessagefrom'evamq/message';constmanager=newMQ(config,console);constproducer=awaitmanager.getProducer();constmsg=awaitproducer.produce(newMessage({foo: 'bar'}));console.log('[%s] produced %s',producer.name,msg.toDebugString());消费消息:
importMQfrom'evamq';constmanager=newMQ(config,console);constconsumer=awaitmanager.getConsumer();awaitconsumer.consuming(async(err,message)=>{console.log('[%s] consuming %j',consumer.name,message);},3);consumer.enableGracefulExit();import{MessageTopic}from'evamq';importMessagefrom'evamq/message';consttopic=newMessageTopic(config,console);awaittopic.factoryKafka();// 或 factoryMns() / factoryNats()constpublisher=topic.getPublisher();awaitpublisher.publish(newMessage({event: 'user.created'}));通过配置切换默认实例:
// config 文件{defaultInstance: 'kafka_default'// 改为// defaultInstance: 'mns_default'}或在运行时手动指定:
awaitmanager.factoryMns('another');constproducer=awaitmanager.getProducer('mns_another');{defaultInstance: 'kafka_default',// 默认实例键名kafka: {default: {connection: {/* Kafka 连接参数 */},defaultQueueName: 'my-queue',},},mns: {default: {connection: {accountId: 'your_account_id',region: 'hangzhou',keyId: 'your_key_id',keySecret: 'your_key_secret',},defaultQueueName: 'my-queue',},},nats: {default: {connection: {/* NATS 连接参数 */},defaultQueueName: 'my-queue',},},}实例键格式为 {adapter}_{configKey},如 kafka_default、mns_another。
App → MessageQueue / MessageTopic
→ Adapter (kafka_adapter | mns_adapter | nats_adapter)
→ peer client(动态加载)
Message / CommandMessage 贯穿 produce / consume
MessageQueue(默认导出)— 队列模式,getProducer/getConsumer为 async,自动初始化适配器MessageTopic— 主题模式,getPublisher/getSubscriber为同步,需先调用factory*注册实例
更多完整示例见 examples/:
| 示例 | 说明 |
|---|---|
manager_producer.ts | 队列模式生产消息(默认实例) |
manager_consumer.ts | 队列模式消费消息 |
manager_switch_connection.ts | 运行时切换 MNS 实例 |
kafka_producer.ts / kafka_consumer.ts | 直接使用 Kafka 适配器 |
mns_producer.ts / mns_consumer.ts | 直接使用 MNS 适配器 |
nats_producer.ts / nats_consumer.ts | NATS JetStream 模式 |
rdkafka_producer.ts / rdkafka_consumer.ts | 底层 rdkafka 封装 |
git clone git@github.com:EvaEngine/EvaQueue.js.git
cd EvaQueue.js
pnpm install
pnpm build
pnpm test- Node.js >= 24.0.0
- pnpm(推荐)或 npm
| 命令 | 说明 |
|---|---|
pnpm build | TypeScript 编译到 lib/ |
pnpm test | 运行测试(node:test + tsx + coverage) |
pnpm lint | ESLint 检查 |
pnpm format | Prettier 格式化 |