引言
在微服务架构盛行的今天,服务间通信的复杂性日益增加。同步HTTP调用带来的耦合、性能瓶颈和故障扩散问题,促使我们转向事件驱动架构。消息队列作为事件驱动架构的核心组件,选型不当会导致系统难以维护。本文将从实战角度出发,对比主流消息队列RabbitMQ和Kafka,并手把手带你实现一个事件驱动的订单处理系统。
消息队列选型:RabbitMQ vs Kafka
核心差异
| 特性 | RabbitMQ | Kafka | |——|———-|——-| | 设计理念 | 消息代理,支持复杂路由 | 分布式日志,高吞吐量 | | 消息模型 | Exchange + Queue | Topic + Partition | | 消息顺序 | 单队列内有序 | 分区内有序 | | 消息持久化 | 支持,但性能较低 | 默认持久化,性能高 | | 消费模式 | Push | Pull | | 典型场景 | 任务调度、异步解耦 | 日志收集、流处理 |
选型建议:
- 需要灵活路由、事务消息、延迟消息 → RabbitMQ
- 高吞吐、日志流、事件溯源 → Kafka
- 简单消息队列 → 两者皆可,但Kafka运维成本较高
实战:订单事件驱动系统
我们将构建一个简单的订单处理系统,包含订单服务、库存服务、通知服务,通过消息队列解耦。
技术栈
- Node.js (模拟服务)
- RabbitMQ (amqplib库)
- Kafka (kafkajs库)
- Docker (本地运行队列)
场景设计
- 用户下单 → 订单服务发布“订单创建”事件
- 库存服务消费事件,扣减库存
- 通知服务消费事件,发送短信
环境准备
# 启动RabbitMQ和Kafka
docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3-management
docker run -d --name kafka -p 9092:9092 apache/kafka:latest
使用RabbitMQ实现
1. 安装依赖
mkdir order-system && cd order-system
npm init -y
npm install amqplib
2. 订单服务:发布事件
// publisher.js
const amqp = require('amqplib');
async function publishOrderCreated(order) {
const connection = await amqp.connect('amqp://localhost');
const channel = await connection.createChannel();
const exchange = 'order.events';
await channel.assertExchange(exchange, 'topic', { durable: true });
const routingKey = 'order.created';
channel.publish(exchange, routingKey, Buffer.from(JSON.stringify(order)));
console.log(`Published order ${order.id}`);
setTimeout(() => { connection.close(); process.exit(0); }, 500);
}
publishOrderCreated({ id: 123, userId: 456, amount: 99.99 });
3. 库存服务:消费事件
// consumer_inventory.js
const amqp = require('amqplib');
async function consume() {
const connection = await amqp.connect('amqp://localhost');
const channel = await connection.createChannel();
const exchange = 'order.events';
await channel.assertExchange(exchange, 'topic', { durable: true });
const queue = await channel.assertQueue('', { exclusive: true });
channel.bindQueue(queue.queue, exchange, 'order.created');
channel.consume(queue.queue, msg => {
const order = JSON.parse(msg.content.toString());
console.log(`Inventory: Deduct stock for order ${order.id}`);
// 模拟库存扣减
channel.ack(msg);
});
}
consume();
4. 通知服务:消费事件
// consumer_notification.js
// 代码类似,只是绑定同一个队列或不同队列,这里省略
注意:RabbitMQ的Topic Exchange支持通配符路由,例如order.*可匹配所有订单事件。
使用Kafka实现
1. 安装依赖
npm install kafkajs
2. 订单服务:生产事件
// producer.js
const { Kafka } = require('kafkajs');
const kafka = new Kafka({ clientId: 'order-service', brokers: ['localhost:9092'] });
const producer = kafka.producer();
async function publishOrderCreated(order) {
await producer.connect();
await producer.send({
topic: 'order-events',
messages: [
{ key: order.id.toString(), value: JSON.stringify(order) },
],
});
console.log(`Published order ${order.id}`);
await producer.disconnect();
}
publishOrderCreated({ id: 123, userId: 456, amount: 99.99 });
3. 库存服务:消费事件
// consumer.js
const { Kafka } = require('kafkajs');
const kafka = new Kafka({ clientId: 'inventory-service', brokers: ['localhost:9092'] });
const consumer = kafka.consumer({ groupId: 'inventory-group' });
async function consume() {
await consumer.connect();
await consumer.subscribe({ topic: 'order-events', fromBeginning: true });
await consumer.run({
eachMessage: async ({ topic, partition, message }) => {
const order = JSON.parse(message.value.toString());
console.log(`Inventory: Deduct stock for order ${order.id}`);
},
});
}
consume();
注意:Kafka的消费者组保证每个分区只被组内一个消费者消费,实现负载均衡。
最佳实践与踩坑记录
1. 消息幂等性
消费端必须实现幂等性,避免重复处理。例如使用订单ID作为唯一键,处理前检查是否已处理。
2. 消息顺序性
- RabbitMQ:单个队列内保证顺序,但多个消费者时需注意。
- Kafka:单个分区内有序,可通过相同key确保进入同一分区。
3. 死信队列
处理失败的消息应转入死信队列,避免阻塞主队列。
4. 监控与告警
- RabbitMQ:Management UI查看队列堆积。
- Kafka:使用Kafka Lag监控消费者偏移量。
总结
本文通过实战对比了RabbitMQ和Kafka,并给出了选型建议。关键在于根据业务场景选择:需要灵活路由选RabbitMQ,需要高吞吐选Kafka。事件驱动架构能有效解耦服务,但需注意幂等性和顺序性。下一步可以探索事件溯源(Event Sourcing)和CQRS模式。