消息队列:这次怎么落地的

如果只能用一句话说消息队列:先把失败复现出来。

这次做消息队列改造,从 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/) 转载或引用必须申明原指尖魔法屋来源及源地址!