消息队列:这次怎么落地的
如果只能用一句话说消息队列:先把失败复现出来。
这次做消息队列改造,从 RabbitMQ 到 Kafka,。
为什么需要消息队列
同步调用的问题
// 同步调用:用户体验差
async function processOrder(order) {
// 1. 扣减库存
await inventoryService.reduceStock(order.items);
// 2. 创建订单
const createdOrder = await orderService.create(order);
// 3. 发送通知
await notificationService.send(createdOrder);
// 4. 更新用户积分
await userService.updatePoints(order.user_id, order.total);
return createdOrder;
}
// 问题:
// 1. 响应慢
// 2. 容易失败
// 3. 难以扩展
// 4. 耦合严重
异步处理
// 异步处理:用户体验好
async function processOrder(order) {
// 1. 扣减库存
await inventoryService.reduceStock(order.items);
// 2. 创建订单
const createdOrder = await orderService.create(order);
// 3. 发送消息到消息队列
await queue.publish('order.created', createdOrder);
return createdOrder;
}
// 其他服务订阅消息
queue.subscribe('order.created', async (order) => {
// 发送通知
await notificationService.send(order);
// 更新用户积分
await userService.updatePoints(order.user_id, order.total);
});
RabbitMQ
基础连接
import pika
# 连接 RabbitMQ
connection = pika.BlockingConnection(
pika.ConnectionParameters('localhost')
)
channel = connection.channel()
# 创建队列
channel.queue_declare(queue='orders')
# 发送消息
def send_order(order):
channel.basic_publish(
exchange='',
routing_key='orders',
body=json.dumps(order)
)
# 接收消息
def callback(ch, method, properties, body):
order = json.loads(body)
process_order(order)
ch.basic_ack(delivery_tag=method.delivery_tag)
channel.basic_consume(
queue='orders',
on_message_callback=callback
)
channel.start_consuming()
交换机和队列
# 声明交换机
channel.exchange_declare(
exchange='orders_exchange',
exchange_type='topic'
)
# 声明队列
channel.queue_declare(queue='orders_queue')
channel.queue_declare(queue='notifications_queue')
# 绑定队列到交换机
channel.queue_bind(
exchange='orders_exchange',
queue='orders_queue',
routing_key='order.created'
)
channel.queue_bind(
exchange='orders_exchange',
queue='notifications_queue',
routing_key='order.*'
)
# 发送消息
def publish_order(order):
channel.basic_publish(
exchange='orders_exchange',
routing_key='order.created',
body=json.dumps(order)
)
死信队列
# 声明死信队列
channel.queue_declare(queue='dlq_orders')
# 声明主队列,配置死信交换机
arguments = {
'x-dead-letter-exchange': '',
'x-dead-letter-routing-key': 'dlq_orders'
}
channel.queue_declare(
queue='orders_queue',
arguments=arguments
)
# 消费者处理失败时,消息会进入死信队列
def callback(ch, method, properties, body):
try:
order = json.loads(body)
process_order(order)
ch.basic_ack(delivery_tag=method.delivery_tag)
except Exception as e:
# 拒绝消息,重新排队或进入死信队列
ch.basic_nack(
delivery_tag=method.delivery_tag,
requeue=False
)
Kafka
基础连接
from kafka import KafkaProducer, KafkaConsumer
import json
# 生产者
producer = KafkaProducer(
bootstrap_servers='localhost:9092',
value_serializer=lambda v: json.dumps(v).encode('utf-8')
)
# 发送消息
def send_order(order):
producer.send('orders', value=order)
producer.flush()
# 消费者
consumer = KafkaConsumer(
'orders',
bootstrap_servers='localhost:9092',
value_deserializer=lambda m: json.loads(m.decode('utf-8'))
)
# 接收消息
for message in consumer:
order = message.value
process_order(order)
分区和副本
from kafka import KafkaAdminClient
from kafka.admin import NewTopic
# 创建主题,指定分区和副本
admin_client = KafkaAdminClient(bootstrap_servers='localhost:9092')
topic = NewTopic(
name='orders',
num_partitions=3,
replication_factor=2
)
admin_client.create_topics([topic])
# 生产者指定分区
def send_order(order):
# 根据 user_id 计算分区
partition = hash(order['user_id']) % 3
producer.send('orders', value=order, partition=partition)
producer.flush()
# 消费者订阅分区
consumer.subscribe(['orders'])
# 或指定订阅特定分区
consumer.assign([TopicPartition('orders', 0)])
消费者组
# 创建消费者组
consumer = KafkaConsumer(
'orders',
bootstrap_servers='localhost:9092',
group_id='order-processing-group',
value_deserializer=lambda m: json.loads(m.decode('utf-8'))
)
# 同一个消费者组的消费者会自动分区消费
# 例如:有 3 个分区,2 个消费者
# 消费者 1 消费分区 0, 1
# 消费者 2 消费分区 2
for message in consumer:
order = message.value
process_order(order)
# 自动提交偏移量
consumer.commit()
消息模式
发布订阅
# RabbitMQ 发布订阅
channel.exchange_declare(exchange='notifications', exchange_type='fanout')
channel.queue_declare(queue='email_notifications')
channel.queue_bind(exchange='notifications', queue='email_notifications')
channel.queue_declare(queue='sms_notifications')
channel.queue_bind(exchange='notifications', queue='sms_notifications')
# 发送消息
def publish_notification(notification):
channel.basic_publish(
exchange='notifications',
routing_key='',
body=json.dumps(notification)
)
# Kafka 发布订阅
producer.send('notifications', value=notification)
# 消费者订阅同一个主题
consumer.subscribe(['notifications'])
请求响应
# 临时队列
response_queue = channel.queue_declare(queue='', exclusive=True)
callback_queue = response_queue.method.queue
# 设置回调队列
channel.basic_consume(
queue=callback_queue,
on_message_callback=lambda ch, method, props, body: handle_response(props.correlation_id, body)
)
# 发送请求
def send_request(request):
correlation_id = str(uuid.uuid4())
channel.basic_publish(
exchange='',
routing_key='rpc_queue',
properties=pika.BasicProperties(
reply_to=callback_queue,
correlation_id=correlation_id
),
body=json.dumps(request)
)
return correlation_id
def handle_response(correlation_id, body):
response = json.loads(body)
print(f"Received response: {response}")
延迟队列
# 使用 RabbitMQ 延迟插件
channel.queue_declare(
queue='delayed_orders',
arguments={'x-delayed-type': 'direct'}
)
channel.exchange_declare(exchange='delayed_exchange', exchange_type='direct')
channel.queue_bind(
exchange='delayed_exchange',
queue='delayed_orders'
)
# 发送延迟消息
def send_delayed_order(order, delay_seconds):
channel.basic_publish(
exchange='delayed_exchange',
routing_key='delayed_orders',
properties=pika.BasicProperties(
headers={'x-delay': delay_seconds * 1000}
),
body=json.dumps(order)
)
踩过的坑
坑一:消息丢失
消息在传输过程中丢失。
解决:使用持久化、确认机制。
# RabbitMQ 持久化
channel.queue_declare(
queue='orders',
durable=True
)
channel.basic_publish(
exchange='',
routing_key='orders',
body=json.dumps(order),
properties=pika.BasicProperties(
delivery_mode=2 # 持久化消息
)
)
# Kafka 持久化
producer.send(
'orders',
value=order,
acks='all' # 等待所有副本确认
)
坑二:消息重复
消息被重复消费。
解决:使用消息去重。
import redis
import hashlib
# 使用 Redis 去重
redis_client = redis.Redis(host='localhost', port=6379, db=0)
def process_order(order):
# 计算消息 ID
message_id = hashlib.md5(json.dumps(order).encode()).hexdigest()
# 检查是否已处理
if redis_client.exists(f"processed:{message_id}"):
print("Message already processed")
return
# 处理消息
process_order_logic(order)
# 标记为已处理
redis_client.setex(
f"processed:{message_id}",
3600, # 1 小时过期
1
)
坑三:消息顺序
消息顺序错乱。
解决:使用有序队列或分区。
# RabbitMQ:单队列保证顺序
channel.queue_declare(queue='orders_queue')
# Kafka:同一分区保证顺序
def send_order(order):
# 根据用户 ID 分区
partition = hash(order['user_id']) % 3
producer.send('orders', value=order, partition=partition)
坑四:消费者积压
消费者处理速度慢,消息积压。
解决:增加消费者或优化处理逻辑。
# 增加消费者
def start_consumers(count=3):
for i in range(count):
thread = threading.Thread(target=consume_messages)
thread.start()
# 优化处理逻辑
def process_order(order):
# 批量处理
batch_size = 100
orders = []
for message in consumer:
orders.append(message.value)
if len(orders) >= batch_size:
process_orders_batch(orders)
orders = []
写在最后
消息队列这东西,不只是技术,是架构和团队协作。
解决了:
- 服务解耦
- 异步处理
- 削峰填谷
- 可靠性
带来了:
- 复杂度增加
- 调试困难
- 运维成本
选型之前先评估:
- 消息量
- 实时性要求
- 可靠性要求
- 团队能力
RabbitMQ 适合:
- 低到中等消息量
- 需要复杂的路由
- 需要消息确认
Kafka 适合:
- 高消息量
- 需要高吞吐量
- 需要持久化和重放
不是所有场景都需要消息队列,有时候同步调用更简单。
这次消息队列改造花了一个月,从 RabbitMQ 到 Kafka。改造完成后,系统响应时间减少了 60%,系统可用性从 99.5% 提升到 99.9%。
版权声明: 本文首发于 指尖魔法屋-消息队列:这次怎么落地的(https://blog.thinkmoon.cn/post/88-message-queue-rabbitmq-kafka-practice/) 转载或引用必须申明原指尖魔法屋来源及源地址!
评论
使用 GitHub 账号登录后即可留言,支持 Markdown。