← 返回资讯
林远舟
技术编辑
已审核

事件驱动架构:从 RabbitMQ 到 Kafka 的实战对比

title: "事件驱动架构:从 RabbitMQ 到 Kafka 的实战对比"

事件驱动架构:从 RabbitMQ 到 Kafka 的实战对比

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(任务分发、通知)

事件驱动架构的核心不是选择哪个消息中间件,而是理解事件作为一等公民的设计思维。先设计好事件模型,再选择合适的技术实现。

402
8045 阅读
3 评论
分享
链接已复制
编辑说明

本文由 MakeSense 编辑团队撰写并审核。文中引用的数据和观点均经过交叉验证,如有疏漏欢迎在评论区指正。最后更新:2026年07月11日 08:50

林远舟

技术编辑

全栈工程师出身,做过 5 年技术社区运营。对 AI 编程工具、开发者生态有深入研究,喜欢用实测数据说话。

读者评论 3

M
创业者Mark 1周前
正在做相关方向,这篇文章给了我不少启发。
回复 点赞 (7)
老李 1周前
有个小问题想请教,文中提到的那个方案在大规模场景下性能怎么样?
回复 点赞 (5)
运营小陈 2天前
转发到团队群了,大家都觉得有参考价值。
回复 点赞 (4)