Skip to content
Published on

[Architecture] 用消息队列实现异步系统:Kafka vs RabbitMQ

分享
Authors

概述

在现代分布式系统中,服务之间的通信正变得越来越复杂。 仅靠同步方式的 HTTP 调用,很难达成高吞吐量、故障隔离与灵活伸缩。 本文整理以 消息队列(Message Queue)为基础的异步架构核心概念, 深入对比 Apache Kafka 与 RabbitMQ,并覆盖 LLM Gateway 中的实战异步处理模式。


1. 同步 vs 异步架构

1.1 同步(Synchronous)通信的局限

在同步通信中,客户端向服务器发送请求后会一直等待响应返回。

Client --> Service A --> Service B --> Service C
          (blocking)    (blocking)    (processing)

问题点:

  • Tight Coupling:只要一个服务宕机,整条链路就会失败
  • Latency 累积:各个服务的响应时间会累加
  • 伸缩受限:最慢的服务决定了整体吞吐量
  • 资源浪费:等待响应期间一直占用线程/连接

1.2 异步(Asynchronous)通信的优势

引入消息队列后,生产者与消费者被解耦。

Producer --> [Message Queue] --> Consumer
  (fire)    (buffer/persist)    (process at own pace)

优点:

  • Loose Coupling:生产者与消费者独立运行
  • Load Leveling:流量激增时队列充当缓冲
  • Fault Isolation:消费者故障不会影响生产者
  • Independent Scaling:消费者可以独立横向扩展
  • Peak Shaving:把瞬时负载平滑化

2. Message Queue 核心概念

2.1 基本组成要素

组成要素作用
Producer生成消息并发送到队列/主题
Consumer从队列/主题接收消息并处理
Broker接收、存储、转发消息的中间服务器
Topic / Queue存放消息的逻辑通道
Partition把主题物理拆分,从而支持并行处理

2.2 消息投递保证级别

At-Most-Once:   Producer --> Broker (no retry)
                消息可能丢失,但不会重复

At-Least-Once:  Producer --> Broker --> ACK --> Retry on failure
                消息不会丢失,但可能重复

Exactly-Once:   Producer --> Broker (idempotent + transaction)
                消息精确投递一次

2.3 消息投递模型

Point-to-Point (Queue):

  • 一条消息只投递给一个消费者
  • 适合任务分发(Work Distribution)

Publish-Subscribe (Topic):

  • 一条消息投递给所有订阅者
  • 适合事件广播

3. Apache Kafka 深入解析

3.1 什么是 Kafka

Apache Kafka 是 LinkedIn 开发的分布式事件流平台。 它以高吞吐量、持久性和水平扩展性为特征,广泛用于实时数据管道与流式应用。

3.2 架构概览

                    +-------------------+
                    |   Kafka Cluster   |
                    |                   |
  Producers ------->|  Broker 1         |-------> Consumers
                    |  Broker 2         |         (Consumer Group A)
                    |  Broker 3         |-------> Consumers
                    |                   |         (Consumer Group B)
                    +-------------------+
                           |
                    +------+------+
                    |  KRaft      |
                    |  Controller |
                    +-------------+

主要组成要素:

  • Broker:接收、存储、转发消息的服务器节点
  • KRaft Controller:从 Kafka 4.0 起替代 ZooKeeper 的内置元数据管理层
  • Topic:消息的逻辑分类
  • Partition:主题的物理分片,每个分区都是保证顺序的不可变日志
  • Replication:按分区维护副本,采用 Leader-Follower 结构
  • ISR (In-Sync Replicas):与 Leader 保持同步的副本集合

3.3 KRaft 模式(移除 ZooKeeper)

从 Kafka 4.0 起,KRaft 成为默认模式。

原有结构:
  Kafka Brokers <---> ZooKeeper Ensemble (独立集群)

KRaft 结构:
  Kafka Brokers (其中一部分兼任 Controller 角色)
  - 通过内部 Raft 共识协议管理元数据
  - 不需要独立的 ZooKeeper 集群

KRaft 的优点:

  • 降低运维复杂度(无需管理 ZooKeeper 集群)
  • 改善分区数量的扩展上限(支持数百万分区)
  • 提升元数据传播速度
  • 缩短控制器故障恢复时间

3.4 Producer 深入解析

Partitioner

决定把消息发送到哪个分区。

from confluent_kafka import Producer

conf = {
    'bootstrap.servers': 'localhost:9092',
    'client.id': 'my-producer',
    'acks': 'all',
    'enable.idempotence': True,
    'max.in.flight.requests.per.connection': 5,
    'retries': 2147483647,
    'linger.ms': 5,
    'batch.size': 16384,
    'compression.type': 'lz4',
}

producer = Producer(conf)

def delivery_callback(err, msg):
    if err:
        print(f"Delivery failed: {err}")
    else:
        print(f"Delivered to {msg.topic()} [{msg.partition()}] @ {msg.offset()}")

# 基于键的分区 - 相同的键进入相同的分区
producer.produce(
    topic='orders',
    key='user-123',
    value='{"order_id": "ord-456", "amount": 5000}',
    callback=delivery_callback
)

producer.flush()

acks 配置

acks 值行为持久性性能
0不等待 broker 响应最高
1仅确认 Leader 写入
all (-1)确认所有 ISR 已写入最高

Batching 与 Compression

Producer 内部动作:
  Record --> Accumulator --> [Batch] --> Compressor --> Network Send

  linger.ms: 批量发送等待时间 (默认 0ms)
  batch.size: 批次最大大小 (默认 16KB)
  compression.type: none / gzip / snappy / lz4 / zstd

3.5 Consumer 深入解析

Consumer Group

Topic: orders (3 partitions)

Consumer Group A:
  Consumer 1 <-- Partition 0
  Consumer 2 <-- Partition 1
  Consumer 3 <-- Partition 2

Consumer Group B:
  Consumer 4 <-- Partition 0, 1, 2

每个 Consumer Group 都独立消费消息。 分区内部保证顺序,并且一个分区在组内只会分配给一个消费者。

Offset 管理

from confluent_kafka import Consumer, KafkaError

conf = {
    'bootstrap.servers': 'localhost:9092',
    'group.id': 'order-processor',
    'auto.offset.reset': 'earliest',
    'enable.auto.commit': False,  # 推荐手动提交
    'max.poll.interval.ms': 300000,
    'session.timeout.ms': 45000,
}

consumer = Consumer(conf)
consumer.subscribe(['orders'])

try:
    while True:
        msg = consumer.poll(timeout=1.0)
        if msg is None:
            continue
        if msg.error():
            if msg.error().code() == KafkaError._PARTITION_EOF:
                continue
            else:
                print(f"Error: {msg.error()}")
                break

        # 处理消息
        process_order(msg.value().decode('utf-8'))

        # 手动提交 - 处理完成之后
        consumer.commit(asynchronous=False)
finally:
    consumer.close()

Rebalancing

消费者组的成员发生变化时,会触发分区的重新分配。

Rebalancing 触发条件:
  1. 新的消费者加入
  2. 已有消费者离开 (crash 或正常退出)
  3. 主题分区数变更
  4. 订阅的主题变更

Rebalancing 策略:
  - Eager (Stop-the-World): 释放全部分区后重新分配
  - Cooperative (Incremental): 只重新分配发生变化的分区

3.6 Exactly-Once Semantics

# Idempotent Producer + Transactional Messaging
conf = {
    'bootstrap.servers': 'localhost:9092',
    'enable.idempotence': True,
    'transactional.id': 'my-transaction-id',
}

producer = Producer(conf)
producer.init_transactions()

try:
    producer.begin_transaction()

    producer.produce('topic-a', key='k1', value='v1')
    producer.produce('topic-b', key='k2', value='v2')

    # Consumer offset 也纳入事务
    producer.send_offsets_to_transaction(
        consumer.position(consumer.assignment()),
        consumer.consumer_group_metadata()
    )

    producer.commit_transaction()
except Exception as e:
    producer.abort_transaction()
    raise e

Exactly-Once 的工作原理:

  1. Idempotent Producer:用 Producer ID + Sequence Number 防止重复发送
  2. Transactional Messaging:对多个主题/分区做原子写入
  3. read_committed isolation:只消费已提交事务的消息

3.7 Kafka Connect 与 Kafka Streams

Kafka Connect:
  Source Connector: DB/File/API --> Kafka Topic
  Sink Connector:   Kafka Topic --> DB/File/API

  示例: Debezium (MySQL CDC) --> Kafka --> Elasticsearch Sink

Kafka Streams:
  Kafka Topic --> Stream Processing --> Kafka Topic

  特点: 不需要独立集群,以库的形式提供
  功能: filter, map, groupBy, windowed aggregation, join

4. RabbitMQ 深入解析

4.1 什么是 RabbitMQ

RabbitMQ 是基于 AMQP(Advanced Message Queuing Protocol)的开源消息代理。 它以灵活的路由、多种协议支持和易于管理为特征。

4.2 架构概览

Producer --> Exchange --> Binding --> Queue --> Consumer

  Exchange: 消息路由中枢
  Binding:  Exchange 与 Queue 之间的路由规则
  Queue:    存放消息的缓冲区
  Virtual Host: 逻辑隔离单位 (多租户)

4.3 Exchange Types

Direct Exchange

Producer --> [Direct Exchange]
              |
              +--(routing_key="order.created")--> [Order Queue]
              |
              +--(routing_key="payment.processed")--> [Payment Queue]
import pika

connection = pika.BlockingConnection(
    pika.ConnectionParameters('localhost')
)
channel = connection.channel()

# 声明 Direct Exchange
channel.exchange_declare(
    exchange='orders',
    exchange_type='direct',
    durable=True
)

# 声明 Queue 并绑定
channel.queue_declare(queue='order_created', durable=True)
channel.queue_bind(
    exchange='orders',
    queue='order_created',
    routing_key='order.created'
)

# 发布消息
channel.basic_publish(
    exchange='orders',
    routing_key='order.created',
    body='{"order_id": "123", "user_id": "456"}',
    properties=pika.BasicProperties(
        delivery_mode=2,  # persistent
        content_type='application/json',
    )
)

Topic Exchange

Producer --> [Topic Exchange]
              |
              +--(routing_key="order.*.kr")--> [Korea Order Queue]
              |
              +--(routing_key="order.premium.#")--> [Premium Queue]
              |
              +--(routing_key="#")--> [Audit Queue] (所有消息)

  * : 精确匹配一个单词
  # : 匹配 0 个以上的单词

Fanout Exchange

Producer --> [Fanout Exchange]
              |
              +--> [Queue A] (所有消息)
              |
              +--> [Queue B] (所有消息)
              |
              +--> [Queue C] (所有消息)

Headers Exchange

不使用路由键,而是用消息头的 key-value 对进行路由。

4.4 Message Acknowledgement 与 Prefetch

def callback(ch, method, properties, body):
    try:
        process_message(body)
        # 处理成功时 ACK
        ch.basic_ack(delivery_tag=method.delivery_tag)
    except Exception as e:
        # 处理失败时 NACK (requeue)
        ch.basic_nack(
            delivery_tag=method.delivery_tag,
            requeue=True
        )

# Prefetch 配置 - 限制一次处理的消息数量
channel.basic_qos(prefetch_count=10)

channel.basic_consume(
    queue='order_created',
    on_message_callback=callback,
    auto_ack=False  # 手动 ACK
)

channel.start_consuming()

Prefetch 策略:

  • prefetch_count=1:逐条处理,分配最为公平
  • prefetch_count=10-50:吞吐量与公平性的平衡
  • prefetch_count=0:不做限制(不推荐)

4.5 Dead Letter Queue (DLQ)

# DLQ 配置
channel.queue_declare(
    queue='order_processing',
    durable=True,
    arguments={
        'x-dead-letter-exchange': 'dlx',
        'x-dead-letter-routing-key': 'order.failed',
        'x-message-ttl': 30000,       # 30 秒 TTL
        'x-max-length': 10000,         # 队列最大长度
    }
)

# DLX (Dead Letter Exchange) 配置
channel.exchange_declare(exchange='dlx', exchange_type='direct')
channel.queue_declare(queue='order_failed', durable=True)
channel.queue_bind(
    exchange='dlx',
    queue='order_failed',
    routing_key='order.failed'
)

进入 Dead Letter 的情况:

  1. Consumer 对消息执行 reject/nack (requeue=False)
  2. 消息 TTL 到期
  3. 超过队列最大长度

4.6 Publisher Confirms

channel.confirm_delivery()

try:
    channel.basic_publish(
        exchange='orders',
        routing_key='order.created',
        body=message_body,
        properties=pika.BasicProperties(delivery_mode=2),
        mandatory=True  # 路由失败时退回
    )
    print("Message confirmed by broker")
except pika.exceptions.UnroutableError:
    print("Message could not be routed")
except pika.exceptions.NackError:
    print("Message was nacked by broker")

4.7 Clustering 与高可用

Classic Mirrored Queue (遗留方案):
  Node 1 (Leader)  <-mirror-> Node 2 (Mirror) <-mirror-> Node 3 (Mirror)
  问题: 存在脑裂风险,同步开销大

Quorum Queue (推荐):
  Node 1 (Leader)  <-Raft-> Node 2 (Follower) <-Raft-> Node 3 (Follower)
  优点: Raft 共识协议,防止脑裂,保证数据安全

Quorum Queue 声明:

channel.queue_declare(
    queue='important_orders',
    durable=True,
    arguments={
        'x-queue-type': 'quorum',
        'x-quorum-initial-group-size': 3,
    }
)

4.8 Streams(Kafka 风格的功能)

RabbitMQ 3.9+ 引入的 Streams 与 Kafka 基于日志的存储相似。

channel.queue_declare(
    queue='order_stream',
    durable=True,
    arguments={
        'x-queue-type': 'stream',
        'x-max-length-bytes': 1073741824,  # 1GB
        'x-stream-max-segment-size-bytes': 52428800,  # 50MB
    }
)

5. LLM Gateway 异步处理模式

5.1 为什么 LLM 调用需要 Queue

LLM API 调用具有与普通 REST API 不同的特性。

特性普通 APILLM API
响应时间50-200ms2-30 秒
成本几乎免费按 token 计费
Rate Limit低(RPM/TPM 限制)
错误率相对较高(429, 500, timeout)
响应大小稳定可变(取决于 token 数)

Queue 解决的问题:

  1. Rate Limit 管理:在队列里控制速率再调用 API
  2. 成本控制:把请求放入队列,按优先级处理
  3. 故障应对:API 故障时先保留在队列中再重试
  4. 负载分散:由 Worker Pool 依次处理大量用户请求

5.2 架构示意图

+----------+     +-------------+     +----------+     +--------------+
|          |     |             |     |          |     |              |
|  Client  +---->+ FastAPI     +---->+ Request  +---->+ Worker Pool  |
|  (Web)   |     | Gateway     |     | Queue    |     | (Celery)     |
|          |     |             |     | (Redis)  |     |              |
+----------+     +------+------+     +----------+     +------+-------+
                        |                                     |
                        |            +----------+             |
                        |            |          |             |
                        +<-----------+ Response +<------------+
                        |            | Store    |
                        |            | (Redis)  |
                        |            +----------+
                        v
                  +----------+
                  |  Status  |
                  |  SSE/WS  |
                  +----------+

5.3 Priority Queue 模式

from enum import IntEnum

class Priority(IntEnum):
    CRITICAL = 0   # 实时聊天机器人响应
    HIGH = 1       # 高级用户
    NORMAL = 2     # 普通用户
    LOW = 3        # 批处理
    BACKGROUND = 4 # 非实时分析

# Celery Task with Priority
from celery import Celery

app = Celery('llm_gateway', broker='redis://localhost:6379/0')

@app.task(bind=True, max_retries=3, default_retry_delay=60)
def call_llm(self, request_id, model, messages, priority):
    try:
        response = litellm.completion(
            model=model,
            messages=messages,
            timeout=30,
        )
        store_response(request_id, response)
        return response
    except Exception as exc:
        # Exponential backoff
        retry_delay = 60 * (2 ** self.request.retries)
        raise self.retry(exc=exc, countdown=retry_delay)

5.4 Retry with Exponential Backoff

import asyncio
import random

async def call_llm_with_retry(
    model: str,
    messages: list,
    max_retries: int = 5,
    base_delay: float = 1.0,
    max_delay: float = 60.0,
):
    for attempt in range(max_retries + 1):
        try:
            response = await litellm.acompletion(
                model=model,
                messages=messages,
                timeout=30,
            )
            return response
        except Exception as e:
            if attempt == max_retries:
                raise

            # Exponential backoff + jitter
            delay = min(base_delay * (2 ** attempt), max_delay)
            jitter = random.uniform(0, delay * 0.1)
            actual_delay = delay + jitter

            print(f"Attempt {attempt + 1} failed: {e}")
            print(f"Retrying in {actual_delay:.1f}s...")
            await asyncio.sleep(actual_delay)

5.5 实战实现:FastAPI + Celery + Redis

# app/main.py
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
from celery.result import AsyncResult
import uuid

app = FastAPI()

class LLMRequest(BaseModel):
    model: str = "gpt-4o"
    messages: list
    priority: int = 2

class LLMResponse(BaseModel):
    request_id: str
    status: str
    result: dict = None

@app.post("/api/v1/completions")
async def create_completion(req: LLMRequest):
    request_id = str(uuid.uuid4())

    # 按优先级发送到 Celery 队列
    task = call_llm.apply_async(
        args=[request_id, req.model, req.messages, req.priority],
        queue=f"llm_priority_{req.priority}",
        priority=req.priority,
    )

    return {
        "request_id": request_id,
        "task_id": task.id,
        "status": "queued",
    }

@app.get("/api/v1/completions/{task_id}")
async def get_completion(task_id: str):
    result = AsyncResult(task_id)

    if result.ready():
        return {
            "status": "completed",
            "result": result.get(),
        }
    elif result.failed():
        return {
            "status": "failed",
            "error": str(result.result),
        }
    else:
        return {
            "status": "processing",
        }
# celery_config.py
from celery import Celery

app = Celery('llm_gateway')

app.conf.update(
    broker_url='redis://localhost:6379/0',
    result_backend='redis://localhost:6379/1',
    task_serializer='json',
    result_serializer='json',
    accept_content=['json'],
    timezone='UTC',
    task_routes={
        'tasks.call_llm': {
            'queue': 'llm_default',
        },
    },
    # 优先级队列配置
    task_queue_max_priority=10,
    task_default_priority=5,
    # Worker 并发限制 (应对 Rate Limit)
    worker_concurrency=4,
    # 任务时间限制
    task_soft_time_limit=120,
    task_time_limit=180,
)
# docker-compose.yml
version: '3.8'
services:
  api:
    build: .
    command: uvicorn app.main:app --host 0.0.0.0 --port 8000
    ports:
      - '8000:8000'
    depends_on:
      - redis
    environment:
      - REDIS_URL=redis://redis:6379

  worker-high:
    build: .
    command: celery -A tasks worker -Q llm_priority_0,llm_priority_1 -c 4
    depends_on:
      - redis

  worker-normal:
    build: .
    command: celery -A tasks worker -Q llm_priority_2 -c 8
    depends_on:
      - redis

  worker-low:
    build: .
    command: celery -A tasks worker -Q llm_priority_3,llm_priority_4 -c 2
    depends_on:
      - redis

  redis:
    image: redis:7-alpine
    ports:
      - '6379:6379'

  flower:
    build: .
    command: celery -A tasks flower --port=5555
    ports:
      - '5555:5555'
    depends_on:
      - redis

6. Kafka vs RabbitMQ 对比

6.1 综合对比表

项目Apache KafkaRabbitMQ
消息模型基于 Pull,Log append基于 Push,Queue
吞吐量数百万 msg/sec数万 msg/sec
延迟~5ms (p99)可达亚毫秒
消息保留按配置周期保留 (默认 7 天)消费后删除
顺序保证分区内保证队列内保证
路由基于 Topic/Partition基于 Exchange/Binding (灵活)
协议自有协议AMQP, MQTT, STOMP
伸缩增加分区实现水平扩展集群化 + Sharding
元数据KRaft (内置)Erlang 分布式系统
流处理内置 Kafka Streams无 (需要外部工具)
运维复杂度中到高
学习曲线
消息大小默认 1MB (可配置)无限制 (实际为数 MB)

6.2 性能基准(参考值)

测试环境: 3 broker/node, 3 replica, 消息大小 1KB

Apache Kafka:
  - 单个分区:   ~50,000 msg/sec
  - 12 个分区:  ~800,000 msg/sec
  - 60 个分区:  ~2,000,000+ msg/sec
  - 延迟 p99: 5ms

RabbitMQ (Quorum Queue):
  - 单个队列:   ~20,000 msg/sec
  - 10 个队列:  ~100,000 msg/sec
  - 延迟 p99: 1ms

6.3 选型指南

应该选择 Kafka 的场景:

  • 大规模事件流(日志、指标、点击流)
  • 实现事件溯源模式
  • 实时数据管道(ETL/ELT)
  • 多个消费者需要独立消费同一条消息
  • 需要消息重放(replay)的场景
  • 需要用 Kafka Streams/ksqlDB 做流处理的场景

应该选择 RabbitMQ 的场景:

  • 需要复杂路由模式的场景(Topic, Headers Exchange)
  • 任务队列(Task Queue)模式
  • 低延迟至关重要的场景
  • 需要支持多种协议的场景(AMQP, MQTT, STOMP)
  • 需要按消息做 TTL、Priority 等精细控制的场景
  • 与既有 AMQP 生态集成

在 LLM Gateway 中该怎么选?

小规模/中规模 (每天 10 万请求以下):
  推荐 RabbitMQ + Celery 组合
  - 配置简单,Priority Queue 支持出色
  - 用 Dead Letter Queue 处理失败很方便

大规模 (每天 100 万请求以上):
  推荐基于 Kafka 的实现
  - 高吞吐量,消息保留使得可以重新处理
  - 用 Consumer Group 实现灵活伸缩
  - 把事件日志用于成本/用量分析

7. 实战监控配置

7.1 Kafka 监控

# Prometheus JMX Exporter 配置 (prometheus-jmx-config.yaml)
rules:
  - pattern: 'kafka.server<type=BrokerTopicMetrics, name=MessagesInPerSec>'
    name: kafka_server_messages_in_total
    type: COUNTER
  - pattern: 'kafka.server<type=BrokerTopicMetrics, name=BytesInPerSec>'
    name: kafka_server_bytes_in_total
    type: COUNTER
  - pattern: 'kafka.consumer<type=consumer-fetch-manager-metrics, client-id=(.+)><>records-lag-max'
    name: kafka_consumer_lag_max
    type: GAUGE

核心监控指标:

  • Consumer Lag:消费者跟上生产者的程度
  • ISR Shrink/Expand:副本状态的变化
  • Under-Replicated Partitions:副本不足的分区
  • Request Latency:请求处理时间

7.2 RabbitMQ 监控

核心指标:
  - Queue Depth: 队列中积压的消息数
  - Consumer Utilization: 消费者利用率
  - Publish/Deliver Rate: 每秒发布/投递的消息数
  - Unacked Messages: 等待 ACK 的消息数
  - Memory/Disk Usage: 资源使用量

RabbitMQ Management UI: http://localhost:15672
Prometheus Plugin: rabbitmq_prometheus (内置)

8. 收尾

基于消息队列的异步架构是现代分布式系统的核心组成部分。

核心整理:

  1. Kafka:大规模事件流、基于日志的存储、用 Consumer Group 灵活扩展
  2. RabbitMQ:灵活路由、低延迟,最适合传统消息队列模式
  3. LLM Gateway:用异步处理实现 Rate Limit 管理、成本控制与故障隔离
  4. 选型标准:依据吞吐量/延迟/路由复杂度/消息保留的必要性来决定

两种技术都拥有成熟的生态,按照项目需求来选择即可。 大规模事件流仍然首选 Kafka,任务分发与灵活路由仍然首选 RabbitMQ。


参考资料