事件驱动架构:Kafka/RabbitMQ/Pulsar选型对比

引言

先讲一个我真实踩过的坑。

三年前我在一家电商公司做订单履约系统。大促前夜,订单服务每秒写入 3 万条消息到一个 RabbitMQ 集群,消费者是库存、积分、风控、通知四个下游服务。结果零点一过,RabbitMQ 的队列积压瞬间冲到千万级,内存告警、磁盘刷盘、连接被打满,最后整个链路雪崩。

事后复盘时我们发现一个残酷的事实:问题不在于 RabbitMQ 不好,而在于我们用错了场景。我们把它当成了一个"高吞吐的日志管道"来用,但 RabbitMQ 的设计哲学其实是"智能代理 + 精确路由"。

这件事让我意识到,消息中间件选型不是比谁的 benchmark 数字更漂亮,而是要理解每个系统在存储模型、投递语义、扩展方式上的根本差异。这篇文章,我会从源码和架构层面把 Kafka、RabbitMQ、Pulsar 这三个主流方案拆开讲清楚,并给出可以直接落地的选型决策树。


核心概念:把消息队列想象成快递系统

在深入之前,先用一个生活化类比把三者的本质差异说透。

假设你要寄快递:

  • RabbitMQ 像顺丰的智能分拣中心。你把包裹交给它,它根据地址(Routing Key)、面单类型(Exchange Type)决定送到哪个网点。每个包裹送出去后,分拣中心就把它"忘掉"了。它关心的是精确投递,不是长期存储。
  • Kafka 像一条高速公路上的传送带。所有包裹按顺序放上传送带,永远不删除(除非超过保留期)。谁想看哪一段就自己去看,传送带本身不关心谁在看。它关心的是吞吐和顺序。
  • Pulsar 像顺丰 + 传送带的混合体。包裹存在一个巨大的中央仓库(BookKeeper),分拣中心(Broker)只是无状态的调度层。谁要看就临时去仓库取,取完打个标记。它关心的是存储与计算分离带来的弹性。

技术定义上:

维度 Kafka RabbitMQ Pulsar
存储模型 分区日志(Partitioned Log) 队列(Queue) 分区日志 + 分层存储
消费模型 Pull Push Pull(兼容 Push 语义)
消息语义 至少一次 / 精确一次 至少一次 / 精确一次(需插件) 至少一次 / 精确一次
顺序保证 分区内有序 队列内有序 分区内有序
存储与计算 耦合 耦合 分离

原理深度分析:从源码看本质

Kafka:日志即一切

Kafka 的核心抽象只有两个:Log 和 Index。打开 Kafka 的源码,Log.scala 里最核心的字段是:

// core/src/main/scala/kafka/log/Log.scala(简化)
class Log(@volatile private var _dir: File,
          @volatile private var config: LogConfig,
          ...) {
  // 每个 Partition 对应一个 Log 对象
  // 底层是多个 LogSegment,每个 Segment 是 .log + .index + .timeindex 三个文件
  @volatile private var segments: ConcurrentNavigableMap[java.lang.Long, LogSegment]
}

写入时,Kafka 只是顺序追加到 .log 文件末尾,同时在 .index 里记录 offset 到物理位置的稀疏映射(默认每 4KB 写一条索引)。这就是 Kafka 高吞吐的根源——顺序写 + 零拷贝(sendfile)。

消费者读取时,Kafka 用 FileRecords 的 slice + transferTo 直接把磁盘数据通过 DMA 送到网卡,全程不经过 JVM 堆。这也是为什么 Kafka 单分区能轻松跑到几十万 TPS。

代价在哪里? Kafka 的"智能"全在客户端。Broker 不记录谁消费到哪了(老版本),而是让消费者自己维护 offset。这带来了灵活性,也带来了 rebalance 风暴。

RabbitMQ:Erlang 的 Actor 模型

RabbitMQ 用 Erlang 写,核心是 gen_server 进程模型。每个 Queue 是一个 Erlang 进程,消息在进程邮箱里流动。看它的核心模块 rabbit_amqqueue_process.erl:

%% src/rabbit_amqqueue_process.erl(简化)
handle_call({deliver, Delivery}, From, State) ->
    %% 每个 Queue 是一个独立进程,串行处理投递
    %% 这就是为什么单队列无法水平扩展
    ...

关键点:一个 Queue 就是一个 Erlang 进程,单线程串行。这意味着单队列的吞吐上限被单个 Erlang 进程的调度能力锁死。RabbitMQ 的扩展方式是"多队列 + 多消费者",而不是 Kafka 那种"单分区 + 多消费者组"。

它的 Exchange 路由逻辑在 rabbit_exchange_type_*.erl 里,Direct、Topic、Fanout、Headers 四种类型本质上是对 Routing Key 的不同匹配算法。Topic 用的是 trie 树匹配,性能是 O(k),k 是 key 的段数。

Pulsar:计算存储分离的架构

Pulsar 的架构是三者里最"现代"的:

graph TD P[Producer] --> B1[Broker - 无状态] B1 --> BK1[BookKeeper Bookie 1] B1 --> BK2[BookKeeper Bookie 2] B1 --> BK3[BookKeeper Bookie 3] C[Consumer] --> B2[Broker - 无状态] B2 --> BK1 B2 --> BK2 B2 --> BK3 Z[ZooKeeper/Etcd] --> B1 Z --> B2

Broker 本身不存数据,只做协议解析和路由。真正的存储在 BookKeeper 的 Bookie 节点上。一条消息写入时,Broker 会并发写多个 Bookie(默认 2 写 1 确认,即 ensemble=3, writeQuorum=2, ackQuorum=2),只要满足写入法定人数就返回 ack。

这种设计的好处是:

  1. Broker 无状态,扩缩容秒级完成,不像 Kafka 扩容要迁移分区数据。
  2. 存储独立扩展,Bookie 可以按容量单独加机器。
  3. 分层存储天然友好,冷数据自动卸载到 S3。

代价是:架构复杂度高,需要额外运维 BookKeeper 集群,且延迟比 Kafka 略高(多一跳网络)。


实战代码:三个可运行的示例

示例 1:Kafka 精确一次消费(幂等 + 事务)

这个示例演示 Kafka 的 EOS(Exactly Once Semantics),关键点在于 isolation.level=read_committed 和事务性生产者。

import org.apache.kafka.clients.consumer.*;
import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;

import java.time.Duration;
import java.util.Collections;
import java.util.Properties;

/**
 * Kafka 精确一次消费示例:
 * 从 input-topic 读取,处理后写入 output-topic,全程在一个事务里。
 * 保证「消费 offset 提交」和「生产消息」原子性。
 */
public class KafkaEosDemo {

    private static final String BOOTSTRAP = "localhost:9092";
    private static final String INPUT_TOPIC = "input-topic";
    private static final String OUTPUT_TOPIC = "output-topic";
    private static final String GROUP_ID = "eos-demo-group";

    public static void main(String[] args) {
        // ---------- 生产者配置 ----------
        Properties producerProps = new Properties();
        producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP);
        producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        // 关键1:开启幂等,配合事务使用
        producerProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
        // 关键2:设置事务ID,同一个事务ID重启后会 abort 未完成事务
        producerProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "tx-eos-demo-1");
        // 关键3:所有 ISR 确认,保证不丢
        producerProps.put(ProducerConfig.ACKS_CONFIG, "all");

        // ---------- 消费者配置 ----------
        Properties consumerProps = new Properties();
        consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP);
        consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, GROUP_ID);
        consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        // 关键4:关闭自动提交,改由事务统一管理
        consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
        // 关键5:只读已提交事务的消息,避免读到未提交数据
        consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
        // 关键6:消费者必须能感知事务,否则无法参与事务
        consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

        KafkaProducer<String, String> producer = new KafkaProducer<>(producerProps);
        KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps);

        // 关键7:初始化事务
        producer.initTransactions();
        consumer.subscribe(Collections.singletonList(INPUT_TOPIC));

        try {
            while (true) {
                ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
                if (records.isEmpty()) continue;

                // 关键8:开启事务
                producer.beginTransaction();
                try {
                    for (ConsumerRecord<String, String> record : records) {
                        String processed = record.value().toUpperCase();
                        producer.send(new ProducerRecord<>(OUTPUT_TOPIC, record.key(), processed));
                    }
                    // 关键9:把消费位移也纳入事务,实现原子提交
                    producer.sendOffsetsToTransaction(
                            currentOffsets(records), consumer.groupMetadata());
                    producer.commitTransaction();
                } catch (Exception e) {
                    // 关键10:任一步失败则回滚,下游看不到任何中间态
                    producer.abortTransaction();
                    throw e;
                }
            }
        } finally {
            producer.close();
            consumer.close();
        }
    }

    /** 提取当前批次的位移,用于纳入事务 */
    private static java.util.Map<org.apache.kafka.common.TopicPartition, OffsetAndMetadata>
    currentOffsets(ConsumerRecords<String, String> records) {
        java.util.Map<org.apache.kafka.common.TopicPartition, OffsetAndMetadata> offsets =
                new java.util.HashMap<>();
        for (org.apache.kafka.common.TopicPartition tp : records.partitions()) {
            long next = records.records(tp).get(records.records(tp).size() - 1).offset() + 1;
            offsets.put(tp, new OffsetAndMetadata(next));
        }
        return offsets;
    }
}

踩坑提醒:事务 ID 必须全局唯一且稳定(同一个业务实例用同一个 ID)。如果多个实例共用同一事务 ID,Kafka 会互相 abort 对方的事务,导致消息丢失。

示例 2:RabbitMQ 死信队列 + 延迟重试

RabbitMQ 原生不支持延迟消息,但可以用 TTL + DLX 实现优雅的重试机制。

import com.rabbitmq.client.*;

import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.TimeoutException;

/**
 * RabbitMQ 延迟重试示例:
 * 业务队列消费失败 -> 投递到延迟队列(带 TTL)-> 到期后通过 DLX 回到业务队列
 * 支持多级延迟(5s / 30s / 5min),避免雪崩式重试。
 */
public class RabbitMqRetryDemo {

    private static final String EXCHANGE = "biz.exchange";
    private static final String BIZ_QUEUE = "biz.queue";
    private static final String BIZ_ROUTING_KEY = "biz.route";

    // 三级延迟队列,TTL 分别为 5s / 30s / 5min
    private static final long[] RETRY_DELAYS = {5_000, 30_000, 300_000};

    public static void main(String[] args) throws IOException, TimeoutException {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");
        factory.setUsername("guest");
        factory.setPassword("guest");
        Connection connection = factory.newConnection();
        Channel channel = connection.createChannel();

        // 1. 声明业务交换机
        channel.exchangeDeclare(EXCHANGE, BuiltinExchangeType.DIRECT, true);
        channel.queueDeclare(BIZ_QUEUE, true, false, false, null);
        channel.queueBind(BIZ_QUEUE, EXCHANGE, BIZ_ROUTING_KEY);

        // 2. 声明三级延迟队列,各自绑定自己的 DLX(延迟队列本身不消费)
        for (int i = 0; i < RETRY_DELAYS.length; i++) {
            String retryQueue = "retry.queue." + i;
            Map<String, Object> retryArgs = new HashMap<>();
            // 消息在队列里停留 TTL 后变成死信
            retryArgs.put("x-message-ttl", RETRY_DELAYS[i]);
            // 死信路由回业务交换机
            retryArgs.put("x-dead-letter-exchange", EXCHANGE);
            retryArgs.put("x-dead-letter-routing-key", BIZ_ROUTING_KEY);
            channel.queueDeclare(retryQueue, true, false, false, retryArgs);
        }

        // 3. 消费业务队列,失败则按重试次数投递到对应延迟队列
        channel.basicConsume(BIZ_QUEUE, false, new DefaultConsumer(channel) {
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope,
                                       AMQP.BasicProperties props, byte[] body) throws IOException {
                // 从 headers 里读取已重试次数
                int retryCount = 0;
                Object rc = props.getHeaders() != null ? props.getHeaders().get("x-retry-count") : null;
                if (rc != null) retryCount = Integer.parseInt(rc.toString());

                try {
                    process(body);
                    channel.basicAck(envelope.getDeliveryTag(), false);
                } catch (Exception e) {
                    if (retryCount >= RETRY_DELAYS.length) {
                        // 超过最大重试次数,进死信队列人工处理
                        channel.basicNack(envelope.getDeliveryTag(), false, false);
                        return;
                    }
                    // 投递到对应级别的延迟队列
                    Map<String, Object> headers = new HashMap<>();
                    headers.put("x-retry-count", retryCount + 1);
                    AMQP.BasicProperties newProps = new AMQP.BasicProperties.Builder()
                            .headers(headers)
                            .deliveryMode(2)  // 持久化
                            .build();
                    channel.basicPublish(EXCHANGE, "retry.queue." + retryCount, newProps, body);
                    channel.basicAck(envelope.getDeliveryTag(), false);
                }
            }
        });

        System.out.println("消费者已启动,等待消息...");
    }

    private static void process(byte[] body) {
        // 模拟业务处理,抛异常触发重试
        System.out.println("处理消息: " + new String(body));
    }
}

关键设计点:延迟队列必须不消费,只做"等待室"。消息过期后由 RabbitMQ 自动转发到 DLX,这是 RabbitMQ 最经典的重试模式。相比 basicNack(requeue=true) 的立即重试,这种模式可以避免重试风暴。

示例 3:Pulsar 多主题订阅 + 精确一次

Pulsar 的 EOS 依赖事务 API,配合 SubscriptionType.Exclusive 可以做到端到端精确一次。

import org.apache.pulsar.client.api.*;
import org.apache.pulsar.client.api.transaction.Transaction;

import java.util.concurrent.TimeUnit;

/**
 * Pulsar 精确一次示例:
 * 从 input 主题消费,处理后写入 output 主题,事务保证原子性。
 * 重点展示 Pulsar 事务 API 与 Kafka 的差异。
 */
public class PulsarEosDemo {

    private static final String SERVICE_URL = "pulsar://localhost:6650";
    private static final String INPUT_TOPIC = "persistent://public/default/input-topic";
    private static final String OUTPUT_TOPIC = "persistent://public/default/output-topic";
    private static final String SUBSCRIPTION = "eos-sub";

    public static void main(String[] args) throws Exception {
        PulsarClient client = PulsarClient.builder()
                .serviceUrl(SERVICE_URL)
                // 关键1:开启事务,超时时间 1 分钟
                .enableTransaction(true)
                .build();

        // ---------- 生产者 ----------
        Producer<String> producer = client.newProducer(Schema.STRING)
                .topic(OUTPUT_TOPIC)
                // 关键2:发送超时,避免事务悬挂
                .sendTimeout(30, TimeUnit.SECONDS)
                .create();

        // ---------- 消费者 ----------
        Consumer<String> consumer = client.newConsumer(Schema.STRING)
                .topic(INPUT_TOPIC)
                .subscriptionName(SUBSCRIPTION)
                // 关键3:Exclusive 订阅才能保证顺序 + 事务
                .subscriptionType(SubscriptionType.Exclusive)
                .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)
                .subscribe();

        while (true) {
            Message<String> msg = consumer.receive(1, TimeUnit.SECONDS);
            if (msg == null) continue;

            // 关键4:开启事务
            Transaction txn = client.newTransaction()
                    .withTransactionTimeout(1, TimeUnit.MINUTES)
                    .build()
                    .get();

            try {
                String processed = msg.getValue().toUpperCase();
                // 关键5:生产消息绑定到事务
                producer.newMessage(txn)
                        .value(processed)
                        .send();

                // 关键6:ack 也绑定到事务
                consumer.acknowledgeAsync(msg.getMessageId(), txn).get();

                // 关键7:提交事务,此时生产 + ack 才真正生效
                txn.commit().get();
            } catch (Exception e) {
                // 回滚后消息会被重新投递
                txn.abort().get();
                throw e;
            }
        }
    }
}

与 Kafka 的差异:Pulsar 事务是显式对象,可以跨多个 topic 组合;Kafka 事务是隐式绑定到 producer 会话的。Pulsar 的粒度更细,但也更容易写出"忘记 commit"的 bug。


方案对比:什么场景选什么

graph TD A[业务需求] --> B{是否需要消息回溯/重放?} B -->|是| C{是否要求超大规模/多租户?} B -->|否| D{是否要求复杂路由?} C -->|是| E[Pulsar] C -->|否| F[Kafka] D -->|是| G[RabbitMQ] D -->|否| H{吞吐量级?} H -->|万级以下| G H -->|十万级以上| F

更细粒度的对比:

场景 推荐 理由
日志采集、埋点、流处理 Kafka 高吞吐、顺序写、生态成熟(Flink/Spark)
订单、支付等业务解耦 RabbitMQ 精确路由、灵活 ACK、延迟队列
任务队列、RPC 异步化 RabbitMQ 单条消息确认模型天然匹配
多租户 SaaS 事件总线 Pulsar 租户隔离、分层存储、geo-replication
需要消息回溯的 CQRS Kafka/Pulsar 日志模型天然支持重放
消息量大但冷数据多 Pulsar 分层存储自动卸载到 S3
超低延迟(<1ms) RabbitMQ 内存队列,无刷盘
跨地域复制 Pulsar 原生 geo-replication,Kafka 需 MirrorMaker

最佳实践与避坑指南

Kafka 的坑

  1. 分区数不是越多越好。每个分区对应一个文件句柄 + 一个副本线程,1000 分区以上 rebalance 会明显变慢。经验值:单 Broker 分区数控制在 2000 以内。
  2. acks=1 是丢数据的元凶。哪怕只丢一条订单,赔偿成本也远高于那点延迟。生产环境请用 acks=all + min.insync.replicas=2。
  3. 消费者 rebalance 风暴。用 CooperativeStickyAssignor 替代 RangeAssignor,配合 max.poll.interval.ms 合理设置,避免"消费慢导致被踢出组"。
  4. 不要用 Kafka 做延迟消息。原生不支持,硬做只能靠时间轮 + 外部存储,不如直接用 RabbitMQ 或 Pulsar。

RabbitMQ 的坑

  1. 单队列是性能天花板。一个队列一个 Erlang 进程,单队列 5 万 TPS 基本到顶。要么拆队列,要么换 Kafka。
  2. basicNack(requeue=true) 会导致无限重试。必须配合重试计数或死信队列,否则一条毒消息能拖垮整个消费者。
  3. 镜像队列在 3.8 之后被 Quorum Queue 取代。新项目直接上 Quorum Queue,基于 Raft,一致性和性能都更好。
  4. 连接数比消息量更致命。每个 Channel 都有内存开销,建议用连接池,单连接多 Channel。

Pulsar 的坑

  1. BookKeeper 运维复杂度高。没有专业 SRE 团队慎入,否则出问题排查成本极高。
  2. 事务性能有损耗。事务消息比普通消息慢 30%~50%,非必要不开事务。
  3. 消费位点管理要小心。SubscriptionType.Shared 不保证顺序,需要顺序时用 Exclusive 或 Failover。
  4. 分层存储配置不当会导致读放大。冷数据读回要经过 S3 -> Bookie -> Broker 三段,延迟可能到秒级,不适合交互式场景。

通用原则

  • 消息幂等是底线。任何 MQ 都可能重复投递,消费端必须做幂等(唯一键 + 去重表 / Redis SETNX)。
  • 监控滞后指标。Kafka 看 consumer_lag,RabbitMQ 看 queue_messages_ready,Pulsar 看 msgBacklog。
  • 压测要压到雪崩点。只测正常流量没意义,要测"消费者挂掉后积压 1 小时再恢复"的场景。
  • Schema 演进要设计。用 Avro/Protobuf + Schema Registry,避免上下游字段不兼容导致的线上事故。

总结

回到开头那个大促雪崩的案例。如果让我重新选型:

  • 订单主链路(要求精确路由 + 延迟重试)→ RabbitMQ,Quorum Queue。
  • 埋点与流处理(要求高吞吐 + 回溯)→ Kafka。
  • 跨地域事件总线(要求多租户 + geo-replication)→ Pulsar。

三者不是替代关系,而是不同抽象层次的工具。Kafka 是"分布式日志",RabbitMQ 是"消息代理",Pulsar 是"流存储平台"。选型的本质是问自己:我的业务需要的是"日志"、"队列"还是"流"?

最后留一个延伸思考:随着云原生和 Serverless 的普及,未来消息中间件的竞争点可能不再是吞吐和延迟,而是弹性伸缩能力和存储成本。Pulsar 的计算存储分离架构在这个维度上有先天优势,但 Kafka 的 KRaft + Tiered Storage 也在快速追赶。下一个五年的格局,值得持续关注。


参考资料

  • Kafka 源码:core/src/main/scala/kafka/log/Log.scala
  • RabbitMQ 源码:src/rabbit_amqqueue_process.erl
  • Pulsar 官方文档:Transaction API 与 BookKeeper 架构
  • 《Kafka: The Definitive Guide》第 2 版
  • 《RabbitMQ in Depth》