title: "事件驱动架构:从 RabbitMQ 到 Kafka 的实战对比"
date: "2026-07-10"
tags: ["事件驱动", "RabbitMQ", "Kafka", "消息队列"]
事件驱动架构:从 RabbitMQ 到 Kafka 的实战对比
事件驱动架构(EDA)是构建松耦合、可扩展系统的核心模式。RabbitMQ 和 Kafka 是两个主流的消息中间件,设计理念截然不同。
核心概念
事件
事件是系统中发生的有意义的事实。
PYTHON
from dataclasses import dataclass
from datetime import datetime
import uuid
@dataclass
class Event:
event_id: str
event_type: str
timestamp: datetime
payload: dict
source: str
def create_event(event_type: str, payload: dict, source: str) -> Event:
return Event(
event_id=str(uuid.uuid4()),
event_type=event_type,
timestamp=datetime.utcnow(),
payload=payload,
source=source
)
# 示例
order_created = create_event(
event_type="order.created",
payload={"order_id": "12345", "user_id": "u001", "total": 299.99},
source="order-service"
)RabbitMQ:消息代理
RabbitMQ 是传统的消息代理,支持复杂的路由和确认机制。
生产者
PYTHON
import pika
import json
connection = pika.BlockingConnection(
pika.ConnectionParameters('localhost')
)
channel = connection.channel()
# 声明交换机
channel.exchange_declare(
exchange='orders',
exchange_type='topic',
durable=True
)
# 发布消息
def publish_order_event(event: Event):
routing_key = f"order.{event.event_type.split('.')[1]}"
channel.basic_publish(
exchange='orders',
routing_key=routing_key,
body=json.dumps({
"event_id": event.event_id,
"event_type": event.event_type,
"timestamp": event.timestamp.isoformat(),
"payload": event.payload
}),
properties=pika.BasicProperties(
delivery_mode=2, # 持久化
content_type='application/json'
)
)
print(f"Published: {event.event_type}")消费者
PYTHON
import pika
import json
connection = pika.BlockingConnection(
pika.ConnectionParameters('localhost')
)
channel = connection.channel()
# 声明队列并绑定
channel.queue_declare(queue='inventory_service', durable=True)
channel.queue_bind(
exchange='orders',
queue='inventory_service',
routing_key='order.created'
)
def callback(ch, method, properties, body):
event = json.loads(body)
print(f"Received: {event['event_type']}")
# 处理事件
process_order_created(event['payload'])
# 确认消息
ch.basic_ack(delivery_tag=method.delivery_tag)
def process_order_created(payload: dict):
order_id = payload['order_id']
# 扣减库存
inventory_service.reserve(order_id, payload['items'])
channel.basic_consume(
queue='inventory_service',
on_message_callback=callback
)
channel.start_consuming()路由模式
PYTHON
# Direct: 精确匹配
channel.exchange_declare(exchange='direct_logs', exchange_type='direct')
channel.queue_bind(exchange='direct_logs', queue='error_logs', routing_key='error')
# Topic: 通配符匹配
channel.exchange_declare(exchange='topic_logs', exchange_type='topic')
channel.queue_bind(exchange='topic_logs', queue='all_logs', routing_key='#')
channel.queue_bind(exchange='topic_logs', queue='order_logs', routing_key='order.*')
# Fanout: 广播
channel.exchange_declare(exchange='broadcast', exchange_type='fanout')
channel.queue_bind(exchange='broadcast', queue='service_a')
channel.queue_bind(exchange='broadcast', queue='service_b')Kafka:事件流平台
Kafka 是分布式事件流平台,支持高吞吐、持久化和流处理。
生产者
PYTHON
from kafka import KafkaProducer
import json
producer = KafkaProducer(
bootstrap_servers=['localhost:9092'],
value_serializer=lambda v: json.dumps(v).encode('utf-8'),
key_serializer=lambda k: k.encode('utf-8') if k else None,
acks='all',
retries=3
)
def publish_order_event(event: Event):
future = producer.send(
'orders',
key=event.payload.get('order_id'),
value={
"event_id": event.event_id,
"event_type": event.event_type,
"timestamp": event.timestamp.isoformat(),
"payload": event.payload
}
)
# 异步回调
future.add_callback(lambda metadata: print(f"Sent to partition {metadata.partition}"))
future.add_errback(lambda error: print(f"Failed: {error}"))
producer.flush()消费者
PYTHON
from kafka import KafkaConsumer
import json
consumer = KafkaConsumer(
'orders',
bootstrap_servers=['localhost:9092'],
group_id='inventory-service',
auto_offset_reset='earliest',
value_deserializer=lambda m: json.loads(m.decode('utf-8')),
enable_auto_commit=False
)
for message in consumer:
event = message.value
print(f"Received: {event['event_type']} (offset: {message.offset})")
process_order_created(event['payload'])
# 手动提交偏移量
consumer.commit()流处理
PYTHON
from kafka import KafkaConsumer, KafkaProducer
import json
consumer = KafkaConsumer(
'orders',
bootstrap_servers=['localhost:9092'],
value_deserializer=lambda m: json.loads(m.decode('utf-8'))
)
producer = KafkaProducer(
bootstrap_servers=['localhost:9092'],
value_serializer=lambda v: json.dumps(v).encode('utf-8')
)
# 实时聚合
order_counts = {}
for message in consumer:
event = message.value
user_id = event['payload']['user_id']
order_counts[user_id] = order_counts.get(user_id, 0) + 1
# 发送到聚合 topic
producer.send('order-stats', {
'user_id': user_id,
'order_count': order_counts[user_id],
'timestamp': event['timestamp']
})对比
| 特性 | RabbitMQ | Kafka |
|------|----------|-------|
| 消息模型 | 队列/发布订阅 | 日志流 |
| 消息保留 | 消费后删除 | 按时间/大小保留 |
| 吞吐量 | 万级/秒 | 百万级/秒 |
| 延迟 | 微秒级 | 毫秒级 |
| 消息顺序 | 单队列保证 | 分区内保证 |
| 回溯消费 | 不支持 | 支持 |
| 路由能力 | 强(多种交换机) | 弱(仅 Topic) |
| 运维复杂度 | 低 | 高 |
选型建议
选择 RabbitMQ
- 需要复杂路由逻辑
- 消息消费后无需保留
- 团队规模小,运维能力有限
- 对延迟要求极高(微秒级)
选择 Kafka
- 高吞吐场景(日志、事件流)
- 需要消息回溯和重放
- 需要流处理能力
- 事件溯源架构
混合使用
CODE
用户操作 → Kafka(事件流、日志)
↓
流处理聚合
↓
RabbitMQ(任务分发、通知)事件驱动架构的核心不是选择哪个消息中间件,而是理解事件作为一等公民的设计思维。先设计好事件模型,再选择合适的技术实现。
读者评论 3