消息队列学习笔记

从问题出发,讲清消息队列"为什么需要、怎么工作、怎么保证可靠、有什么坑"。不写安装教程,只讲原理和取舍。

📌 学习路线建议

这篇笔记内容较多,不要试图一次全部消化。建议分四步走:

阶段 读哪些章节 目标
第一步:建立直觉 原理篇 第 1~3 节 理解"为什么需要 MQ"、核心模型、四个主流 MQ 的区别
第二步:理解可靠性 原理篇 第 4 节 搞懂消息会丢在哪、怎么幂等、三种语义、顺序保证
第三步:动手跑通 实战篇 第 5~6 节 跑通 Kafka/RabbitMQ 最小示例,理解订单系统完整链路
第四步:生产加固 实战篇 第 7~8 节 死信队列、优雅退出、积压监控、延迟消息、日志收集

每步读完后再进入下一步,不要跳。第三步建议本地 Docker 起一个 Kafka 或 RabbitMQ,把代码跑起来——看十遍不如跑一遍。


第一部分 · 原理篇

这部分讲"为什么"和"怎么想",结合少量关键代码片段帮助理解。读完这部分你能搞懂消息队列的全貌:为什么需要、怎么工作、怎么保证可靠、有什么坑。完整工程代码见第二部分·实战篇。


1 · 为什么需要消息队列

1.1 问题:同步调用的三大痛点

假设你做了一个电商系统,用户下单后要:扣库存 → 生成订单 → 发短信 → 发邮件 → 更新推荐系统。全部同步调用:

  • 慢——每一步都要等上一步完成,用户看着转圈圈
  • 耦合——发短信挂了,订单也跟着失败
  • 扛不住峰值——双 11 瞬间 10 万订单,数据库直接打满

1.2 解法:加个"信箱"

在订单系统和下游服务之间放一个"信箱"(消息队列):

  • 订单系统把消息丢进信箱,立即返回"下单成功"——异步化
  • 短信、邮件、推荐各自从信箱取消息,互不影响——解耦
  • 10 万订单涌入时,信箱先存着,下游按自己的速度消费——削峰

📌 说白了就是生产者-消费者模型。你在操作系统课里学过的经典模式:一个线程往队列里塞数据(生产者),另一个线程从队列里取数据(消费者)。消息队列就是这个模式的工程化升级——在分布式环境下加了持久化、多消费者、路由、重试、监控等基础设施。核心思想没变:用一个中间缓冲区,解耦生产者和消费者的速度差异。

1.3 三个核心价值

价值 同步调用 加消息队列
异步化 用户等所有步骤完成才返回 丢进队列立即返回,下游慢慢处理
解耦 上游知道下游每个服务的地址和接口 上游只管丢消息,不关心谁消费
削峰填谷 峰值流量直接打到数据库 队列缓冲,下游按固定速率消费

2 · 消息队列的核心模型

2.1 两种通信模式

  • 点对点(P2P)——一条消息只能被一个消费者消费。像取快递:取走就没了。
  • 发布/订阅(Pub/Sub)——一条消息可以被多个订阅者各消费一次。像广播:所有人都能听到。

2.2 两种消费方式

  • Push(推)——Broker 主动推给消费者。延迟低,但消费者可能被压垮。
  • Pull(拉)——消费者主动拉取。可控速率,但有空轮询问题。

2.3 核心概念速查

概念 类比 说明
Producer 寄信人 发消息的一方
Consumer 收信人 收消息的一方
Broker 邮局 存储和转发消息的中间件
Topic / Queue 信箱地址 消息的分类
Partition 信箱分格 同一 Topic 的分片,用于并行消费
Offset 取件码 消费者在分区中的位置标记
Consumer Group 合伙取件 同组内每个消费者各消费一部分分区

3 · 主流消息队列对比

3.0 它们分别是什么

  • Kafka——LinkedIn 开源,Apache 顶级项目。和 Nginx 一样,是一个需要独立部署运行的中间件软件。

    什么是"以日志数据结构为核心"?

    这里的"日志(Log)"不是指运行日志(比如 Nginx 的 access.log),而是一种数据结构。你可以把它想象成:

    • 一个只能往末尾追加、不能修改、不能删除中间内容的数组
    • 每条消息进来,就追加到末尾,分配一个递增的编号(Offset)
    • 消费者读消息,就是拿着 Offset 去"翻到第几页看",看完消息还在

    对比传统队列(RabbitMQ):消息像快递,取走就没了。Kafka 的消息像图书馆的书,借阅完放回去,下个人还能借。

    这种"只追加"的设计极其简单——没有删除、没有修改、没有随机写入,磁盘顺序写速度接近内存,所以吞吐能到百万/秒。消息按 Topic 分类、Partition 分片、Offset 定位,消费完消息还在,可以反复回放。适合大流量场景。

    消息存哪里?会撑爆磁盘吗?

    消息存在磁盘上(不在内存),每个 Partition 对应一组磁盘文件。但不会一直堆积——Kafka 用两种机制自动清理:

    • 按时间删除:默认保留 7 天(retention.ms=604800000),超过自动删
    • 按大小删除:默认每个 Partition 保留 1GB(retention.bytes=1073741824),超过删最旧的
    # 配置 Topic 保留策略:只保留最近 1 小时的消息
    kafka-configs.sh --alter --topic orders \
      --add-config retention.ms=3600000
    
    # 配置:每个 Partition 最多保留 500MB
    kafka-configs.sh --alter --topic orders \
      --add-config retention.bytes=524288000

    所以 Kafka 不是"永远不删",而是"保留一段时间后自动删最旧的"。就像图书馆定期清理过期杂志——不是永久保存,但也不是看完就扔,而是放一段时间让有需要的人翻阅。你可以根据业务需求调整保留时间:日志收集可能只留 3 天,订单事件可能留 30 天。

  • RabbitMQ——基于 AMQP 协议的传统消息队列,同样是独立部署的中间件。核心是 Exchange + Queue 的路由模型:生产者发消息到 Exchange,Exchange 根据路由规则分发到不同 Queue,消费者从 Queue 取。消息"取走就没了",支持复杂的路由规则、优先级、TTL。适合业务解耦、延迟敏感的场景。

  • RocketMQ——阿里开源,Apache 顶级项目。借鉴了 Kafka 的存储模型,但加了事务消息支持(先发半消息 → 执行本地事务 → 提交/回滚),适合电商交易这种"消息和数据库必须一致"的场景。

  • Redis Streams——Redis 5.0 内置的轻量队列。不是独立的中间件,就是 Redis 里的一个数据结构,用 XADD/XREAD 操作。适合已经在用 Redis、消息量不大的小项目。

📌 一句话区分:Kafka 是"高速传送带"(大流量、可回放),RabbitMQ 是"快递员"(灵活路由、取走即删),RocketMQ 是"带事务的 Kafka"(消息和 DB 一致),Redis Streams 是"顺手用的轻量队列"。

它们不是某种语言的库,而是独立运行的服务

Kafka、RabbitMQ、RocketMQ 本身是用什么语言写的无所谓——就像 MySQL 是用 C/C++ 写的,但你用 Java、Python、Go、C++ 都能连它。它们都是独立部署运行的服务进程,你的程序通过网络协议连上去读写消息。

中间件 服务端语言 你的程序怎么连 C++ 客户端
Kafka Scala/Java TCP(自定义协议) cppkafka / librdkafka
RabbitMQ Erlang AMQP 协议(TCP) amqpcpp
RocketMQ Java TCP(自定义协议) rocketmq-client-cpp
Redis Streams C RESP 协议(TCP) redis-plus-plus / hiredis

所以本文的 C++ 代码示例,换成 Java、Python、Go 也完全一样——只是换个客户端库,中间件还是那个中间件。

3.1 对比表

Kafka RabbitMQ RocketMQ Redis Streams
定位 日志/流处理 业务消息路由 事务消息 轻量级队列
吞吐 极高(百万/s) 中(万/s) 高(十万/s) 中
延迟 毫秒级 微秒级 毫秒级 微秒级
顺序保证 分区内有序 队列内有序 分区内有序 不保证
持久化 磁盘日志 可选持久化 磁盘 可选 AOF
事务消息 不支持 不支持 支持 不支持
适合场景 日志收集、流计算 复杂路由、低延迟 电商交易、订单 小项目、已有 Redis

📌 选型直觉:

  • 日志/数据管道 → Kafka
  • 复杂路由、业务解耦 → RabbitMQ
  • 交易/订单、需要事务消息 → RocketMQ
  • 已经在用 Redis、量不大 → Redis Streams

4 · 可靠性原理:消息会丢吗?会重复吗?

这一节讲"什么环节会出问题、用什么机制兜底",配合关键代码片段。完整工程代码见第二部分·实战篇。

4.1 消息丢失的三个环节

图 1 \xb7 消息丢失的三个环节

消息从生产到消费经过三个环节,每个环节都可能丢消息:

① 发送环节:消息没到 Broker

问题:网络抖动,消息发出去了但 Broker 没收到。

解法:不要"发了就不管",要等 Broker 明确说"我收到了"。

// Kafka:acks=all → 等所有副本写入才算成功 + 自动重试
Configuration config = {
    {"acks", "all"},        // ★ 关键:等所有副本确认
    {"retries", 3},         // ★ 网络抖动时自动重试
};
Producer producer(config);
producer.produce(MessageBuilder("orders").payload(body));
producer.flush();  // ★ 阻塞等待确认,不是"发完就不管"
// RabbitMQ:开启 confirm 模式,Broker 收到后回调确认
channel.confirm()
    .onSuccess([]() { /* Broker 确认收到 */ })
    .onError([](const char* msg) { /* 发送失败 */ });
channel.publish("orders", "order.created", body,
    AMQP::Mandatory,  // ★ 路由失败时回调,不会静默丢弃
    AMQP::Message().setDeliveryMode(AMQP::DeliveryMode::persistent)  // ★ 持久化
);

② 存储环节:Broker 收到了但没落盘就宕机

问题:Broker 收到消息存在内存里,还没写磁盘就宕机了。

解法:消息要落盘 + 多副本,不能只存内存。

// RabbitMQ:声明持久化队列
channel.declareQueue("orders", AMQP::durable);  // ★ durable:队列元数据也持久化
# Kafka:创建 Topic 时指定 3 副本 + 3 分区
kafka-topics.sh --create --topic orders \
  --partitions 3 --replication-factor 3  # ★ 挂 2 个都不丢

③ 消费环节:消费者拿到消息还没处理完就挂了

问题:消费者刚拿到消息,处理到一半进程崩了。Broker 以为已经消费了,消息永久丢失。

解法:先处理完再确认,不要"拿到就确认"。

// Kafka:关闭自动提交,处理完手动提交 offset
Configuration config = {
    {"enable.auto.commit", false},  // ★ 关键:关掉自动提交
};
Consumer consumer(config);
consumer.subscribe({"orders"});

while (true)
{
    Message msg = consumer.poll();
    if (!msg) continue;
    json order = json::parse(msg.get_payload());
    try {
        process_order(order);    // 先处理业务
        consumer.commit(msg);    // ★ 成功后才提交 offset
    } catch (const std::exception& e) {
        // 不 commit,下次会重新消费这条消息
    }
}
// RabbitMQ:手动 ACK,处理完再确认
channel.consume("orders", AMQP::noack)
    .onReceived([](const AMQP::Message& msg, uint64_t tag, bool redelivered)
    {
        try {
            process_order(order);
            channel.ack(tag);      // ★ 成功 → 确认
        } catch (...) {
            channel.nack(tag, false);  // ★ 失败 → 拒绝,转死信队列
        }
    });

4.2 消息重复:为什么难以避免

网络不可靠 → 重试 → 同一条消息可能被发送/消费多次。根本原因:ACK 可能丢失,发送方无法区分"对方没收到"和"对方收到了但 ACK 丢了"。

解法:消费端幂等——同一消息处理多次,结果和只处理一次一样。两种主流方案:

方案 做法 优点 缺点
Redis SETNX 用消息 ID 做 key,SETNX 标记已处理,设 24h 过期 简单、快 依赖 Redis,Redis 挂了可能重复
数据库唯一约束 建一张已处理消息表,用消息 ID 做主键,INSERT 失败即重复 可靠、和业务同事务 多一次 DB 写入
// 方案一:Redis SETNX 幂等
bool first_time = redis.set("processed:" + order_id, "1",
    SetOptions{SetOptions::SetType::NX, std::chrono::seconds(86400)});
if (!first_time) return;  // ★ 已处理过,跳过
// 执行业务逻辑——只会执行一次
-- 方案二:数据库唯一约束去重
CREATE TABLE processed_messages (
    message_id VARCHAR(64) PRIMARY KEY,  -- 消息唯一 ID
    processed_at TIMESTAMP DEFAULT NOW()
);

📌 核心思想:不要试图阻止重复(网络不可控),而是让重复无害(幂等)。

4.3 三种语义

语义 含义 难度
At Most Once 最多一次,可能丢 最简单(发完不管,acks=0)
At Least Once 至少一次,可能重复 中等(acks=all + 手动提交,大多数 MQ 默认)
Exactly Once 恰好一次,不丢不重 最难(需要事务或幂等保证)

📌 真相:真正的 Exactly Once 几乎不存在,通常是 "At Least Once + 幂等" 伪装出来的。Kafka 的事务(transactional.id)能做到端到端 Exactly Once,但只限 Kafka 内部(Producer→Topic→Consumer),一旦消费侧写外部系统(数据库),还是得靠幂等。

4.4 顺序保证

  • 全局有序——所有消息严格按顺序,代价是只能一个分区、一个消费者,吞吐极低
  • 分区有序——同一 Key 的消息发到同一分区,分区内有序。大多数场景够用
// Kafka:用 order_id 做 key,同一订单的消息一定进同一分区,分区内有序
producer.produce(MessageBuilder("orders")
    .key(order_id)       // ★ key 相同 → 同一 partition → 消费时按 offset 顺序
    .payload(order_data)
);

第二部分 · 实战篇

这部分是完整的工程代码:从快速上手到订单系统、生产加固、更多场景。建议本地起一个 Kafka 或 RabbitMQ,把代码跑起来。


5 · 快速上手:Kafka vs RabbitMQ

5.1 Kafka:最小可运行示例

安装(Docker 一键起):

# docker-compose.yml
services:
  kafka:
    image: bitnami/kafka:3.7
    ports:
      - "9092:9092"
    environment:
      KAFKA_CFG_NODE_ID: 0
      KAFKA_CFG_PROCESS_ROLES: controller,broker
      KAFKA_CFG_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
      KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 0@kafka:9093
      KAFKA_CFG_CONTROLLER_LISTENER_NAMES: CONTROLLER
      KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
docker compose up -d
# C++ 用 cppkafka(librdkafka 的 C++ 封装)
# vcpkg install cppkafka nlohmann-json

生产者 + 消费者:

#include <cppkafka/cppkafka.h>
#include <nlohmann/json.hpp>
#include <iostream>
using namespace cppkafka;
using json = nlohmann::json;

// --- 生产者 ---
Producer producer({
    {"metadata.broker.list", "localhost:9092"}
});

json msg = {{"msg", "hello kafka"}};
producer.produce(MessageBuilder("test").payload(msg.dump()));
producer.flush();  // 确保消息发出去
std::cout << "已发送" << std::endl;

// --- 消费者 ---
Consumer consumer({
    {"metadata.broker.list", "localhost:9092"},
    {"group.id", "test-group"},
    {"auto.offset.reset", "earliest"},
});
consumer.subscribe({"test"});

while (true)
{
    Message msg = consumer.poll();
    if (msg)
    {
        json data = json::parse(msg.get_payload());
        std::cout << "收到: " << data << std::endl;
        break;  // 只收一条
    }
}

5.2 RabbitMQ:最小可运行示例

安装:

docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3-management
# 管理界面:http://localhost:15672  账号/密码:guest/guest
# vcpkg install amqpcpp nlohmann-json

生产者 + 消费者:

#include <amqpcpp.h>
#include <nlohmann/json.hpp>
#include <iostream>
using json = nlohmann::json;

// 假设已有 AMQP::Connection connection 和 AMQP::Channel channel(connection)

// --- 生产者 ---
channel.declareQueue("test");  // 声明队列

json payload = {{"msg", "hello rabbitmq"}};
channel.publish("", "test", payload.dump());
std::cout << "已发送" << std::endl;

// --- 消费者 ---
channel.consume("test", AMQP::noack)  // ★ noack = 手动确认
    .onReceived([](const AMQP::Message& msg, uint64_t tag, bool redelivered)
    {
        json data = json::parse(msg.body());
        std::cout << "收到: " << data << std::endl;
        channel.ack(tag);  // 手动确认
    });
std::cout << "等待消息..." << std::endl;
// channel.startConsuming();  // 启动事件循环

5.3 C++ 版 Kafka 生产者:可靠性配置详解

5.1 的最小示例只配了 metadata.broker.list,生产环境远远不够。下面逐行解释每个关键配置参数:

#include <cppkafka/cppkafka.h>
#include <nlohmann/json.hpp>
#include <iostream>
using namespace cppkafka;
using json = nlohmann::json;

int main()
{
    Configuration config = {
        {"metadata.broker.list", "localhost:9092"},

        // ★ 以下三个参数是生产环境必须配置的可靠性三件套

        {"acks", "all"},
        // acks 控制多少副本确认才算"发送成功"
        //   "0"    → 不等确认,发完就返回(最快,但可能丢消息)
        //   "1"    → 只等 Leader 副本确认(Leader 挂了仍可能丢)
        //   "all"  → 等所有 ISR 副本确认(最安全,Leader 挂了也不丢)

        {"retries", 3},
        // 网络抖动时自动重试次数。配合 retry.backoff.ms 做退避
        //   设 3 次:第 1 次失败 → 等 500ms → 第 2 次 → 等 1000ms → 第 3 次

        {"retry.backoff.ms", 500},
        // 重试间隔,避免狂打 Broker。每次重试间隔翻倍(指数退避)
    };

    Producer producer(config);

    // 构造消息体:用 nlohmann/json 序列化为 JSON 字符串
    json order = {{"order_id", "A001"}, {"amount", 99.9}};
    std::string data = order.dump();

    // ★ produce() 是异步的——消息放进本地缓冲队列,立刻返回
    producer.produce(MessageBuilder("orders").payload(data));

    // ★ flush() 是同步的——阻塞等待所有缓冲消息被 Broker 确认
    //   不调 flush,程序可能消息还没发出去就退出了
    //   生产环境通常在程序退出前或定时调用 flush
    producer.flush();

    std::cout << "发送成功" << std::endl;
    return 0;
}

📌 关键对比:5.1 的最小示例 vs 5.3 的可靠性配置

配置项 5.1 最小示例 5.3 可靠性配置 区别
acks 未设置(默认 1) all 默认只等 Leader,all 等所有副本
retries 未设置(默认 0) 3 默认不重试,设 3 次网络抖动自动恢复
retry.backoff.ms 未设置 500 退避间隔,避免重试风暴
flush() 调了 调了 相同——都要确保消息真正发出

一句话:最小示例能跑通,但生产环境必须加 acks=all + retries + 退避,这是第 4 节"发送环节"原理的代码落地。


6 · 完整实战:订单异步处理系统

把前面学的全部串起来。场景:电商下单,主流程只做"创建订单",后续的扣库存、发短信、加积分全部异步消费。

6.1 架构

图 2 \xb7 订单异步处理架构

6.2 订单服务(生产者)

#include <cppkafka/cppkafka.h>
#include <nlohmann/json.hpp>
#include <mysqlx/xdevapi.h>
#include <boost/uuid/uuid.hpp>
#include <boost/uuid/uuid_generators.hpp>
#include <boost/uuid/uuid_io.hpp>
using namespace cppkafka;
using json = nlohmann::json;

Producer producer({
    {"metadata.broker.list", "localhost:9092"},
    {"acks", "all"},
    {"retries", 3},
});

std::string create_order(const std::string& user_id,
                         const std::string& product_id,
                         int quantity, double amount)
                         {
    // 生成 UUID 作为订单号
    boost::uuids::uuid uuid = boost::uuids::random_generator()();
    std::string order_id = boost::uuids::to_string(uuid);

    // 1. 写数据库——先落盘,这是"事实"
    mysqlx::Session session("localhost", 3306, "root", "123456", "shop");
    auto schema = session.getSchema("shop");
    auto orders = schema.getTable("orders");
    orders.insert("order_id", "user_id", "product_id", "quantity", "amount", "status")
          .values(order_id, user_id, product_id, quantity, amount, "created")
          .execute();
    session.close();

    // 2. 发消息——通知下游处理
    //    ★ 如果发消息失败怎么办?见下面的"本地事务表"方案
    json payload = {
        {"order_id",   order_id},
        {"user_id",    user_id},
        {"product_id", product_id},
        {"quantity",   quantity},
        {"amount",     amount},
    };
    producer.produce(MessageBuilder("order.created")
        .key(order_id)       // ★ 用 order_id 做 key,保证同一订单进同一分区
        .payload(payload.dump())
    );
    producer.flush();

    return order_id;  // 用户立即拿到订单号
}

6.3 库存服务(消费者,幂等 + 手动提交)

#include <cppkafka/cppkafka.h>
#include <nlohmann/json.hpp>
#include <sw/redis++/redis++.h>
#include <mysqlx/xdevapi.h>
using namespace cppkafka;
using json = nlohmann::json;
using namespace sw::redis;

Redis redis("tcp://localhost:6379");

Consumer consumer({
    {"metadata.broker.list", "localhost:9092"},
    {"group.id", "stock-service"},   // ★ 不同的 group:各消费一份完整消息
    {"enable.auto.commit", false},    // ★ 手动提交
    {"auto.offset.reset", "earliest"},
});
consumer.subscribe({"order.created"});

void deduct_stock(const json& order)
{
    std::string order_id = order["order_id"];

    // ★ 幂等检查:Redis SETNX
    std::string key = "stock_done:" + order_id;
    bool first_time = redis.set(key, "1", SetOptions{
        SetOptions::SetType::NX,
        std::chrono::seconds(86400)
    });
    if (!first_time)
    {
        std::cout << "[库存] 订单 " << order_id << " 已扣过库存,跳过" << std::endl;
        return;
    }

    // ★ 数据库层面扣库存
    mysqlx::Session session("localhost", 3306, "root", "123456", "shop");
    auto schema = session.getSchema("shop");
    auto stock = schema.getTable("stock");
    auto result = stock.update()
        .set("quantity", mysqlx::expr("quantity - {}", order["quantity"].get<int>()))
        .where("product_id = :pid AND quantity >= :qty")
        .bind("pid", order["product_id"].get<std::string>())
        .bind("qty", order["quantity"].get<int>())
        .execute();

    if (result.getAffectedItemsCount() == 0)
    {
        std::cerr << "[库存] 商品 " << order["product_id"] << " 库存不足!" << std::endl;
        // 实际项目:发一条"库存不足"消息到死信队列
    } else
    {
        std::cout << "[库存] 订单 " << order_id << " 扣库存成功" << std::endl;
    }
    session.close();
}

while (true)
{
    Message msg = consumer.poll();
    if (!msg || msg.get_error()) continue;

    json order = json::parse(msg.get_payload());
    try {
        deduct_stock(order);
        consumer.commit(msg);  // ★ 处理成功才提交
    } catch (const std::exception& e)
    {
        std::cerr << "[库存] 处理失败: " << e.what() << ",不提交,下次重试" << std::endl;
    }
}

6.4 短信服务(另一个消费者,独立消费)

Consumer sms_consumer({
    {"metadata.broker.list", "localhost:9092"},
    {"group.id", "sms-service"},  // ★ 和库存服务不同的 group,各自独立消费全量消息
    {"enable.auto.commit", false},
    {"auto.offset.reset", "earliest"},
});
sms_consumer.subscribe({"order.created"});

void send_sms(const json& order)
{
    // 实际项目:调用阿里云/腾讯云短信 API
    std::cout << "[短信] 订单 " << order["order_id"]
              << " 已发送短信通知用户 " << order["user_id"] << std::endl;
}

while (true)
{
    Message msg = sms_consumer.poll();
    if (!msg || msg.get_error()) continue;

    json order = json::parse(msg.get_payload());
    try {
        send_sms(order);
        sms_consumer.commit(msg);
    } catch (const std::exception& e)
    {
        std::cerr << "[短信] 发送失败: " << e.what() << ",不提交,下次重试" << std::endl;
    }
}

6.5 消息和数据库的一致性问题

上面的 create_order 有个隐患:数据库写成功了,但发消息失败了怎么办?用户看到"下单成功",但库存永远不扣。

解法:本地事务表(Outbox Pattern)

std::string create_order_safe(const std::string& user_id,
                              const std::string& product_id,
                              int quantity, double amount)
                              {
    // Outbox Pattern:消息和业务数据在同一个事务里写
    boost::uuids::uuid uuid1 = boost::uuids::random_generator()();
    std::string order_id = boost::uuids::to_string(uuid1);
    boost::uuids::uuid uuid2 = boost::uuids::random_generator()();
    std::string message_id = boost::uuids::to_string(uuid2);

    mysqlx::Session session("localhost", 3306, "root", "123456", "shop");

    try {
        session.startTransaction();  // ★ 开启事务
        auto schema = session.getSchema("shop");

        // 1. 写订单
        schema.getTable("orders").insert(
            "order_id", "user_id", "product_id", "quantity", "amount", "status"
        ).values(order_id, user_id, product_id, quantity, amount, "created").execute();

        // 2. 写消息表(和订单在同一个事务里)
        json payload = {
            {"order_id",   order_id},
            {"user_id",    user_id},
            {"product_id", product_id},
            {"quantity",   quantity},
            {"amount",     amount},
        };
        schema.getTable("outbox_messages").insert(
            "message_id", "topic", "payload", "status"
        ).values(message_id, "order.created", payload.dump(), "pending").execute();

        session.commit();  // ★ 两步一起提交,要么都成功要么都失败
    } catch (...)
    {
        session.rollback();
        throw;
    }
    session.close();

    // 3. 异步:一个独立线程/进程扫描 outbox 表,把 pending 的消息发到 Kafka
    //    发送成功后更新 status='sent'
    //    ★ 即使这里挂了也没关系——重启后扫描 outbox 重新发送
    return order_id;
}
// Outbox 发送器(独立线程,定时扫描)
#include <thread>
#include <chrono>

void outbox_sender(Producer& producer)
{
    while (true)
    {
        mysqlx::Session session("localhost", 3306, "root", "123456", "shop");
        auto schema = session.getSchema("shop");
        auto outbox = schema.getTable("outbox_messages");

        // 查询 pending 消息
        auto rows = outbox.select("message_id", "topic", "payload")
            .where("status = 'pending'").limit(100).execute();

        for (auto row : rows)
        {
            std::string message_id = row[0].get<std::string>();
            std::string topic = row[1].get<std::string>();
            std::string payload = row[2].get<std::string>();

            try {
                producer.produce(MessageBuilder(topic).payload(payload));
                producer.flush();
                // 发送成功,更新状态
                outbox.update()
                    .set("status", "sent")
                    .where("message_id = :mid")
                    .bind("mid", message_id)
                    .execute();
            } catch (const std::exception& e)
            {
                std::cerr << "发送 " << message_id << " 失败: " << e.what() << ",下次重试" << std::endl;
            }
        }
        session.close();
        std::this_thread::sleep_for(std::chrono::seconds(1));  // 每秒扫描一次
    }
}

📌 Outbox Pattern 的核心思想:不直接发消息,而是把消息写到数据库的 outbox 表里(和业务数据同一个事务),再由独立进程异步扫描发送。这样保证了数据库写入和消息发送的最终一致性。


7 · 常见坑与最佳实践

7.1 消息积压

原因:消费者处理太慢 / 消费者挂了 / 突发流量 应对:扩消费者、临时扩分区、降级非核心消费

// ★ 监控积压:Kafka 消费延迟 = LogEndOffset - ConsumerOffset
#include <cppkafka/cppkafka.h>
using namespace cppkafka;

long get_lag(const std::string& topic, const std::string& group_id)
{
    Consumer consumer({
        {"metadata.broker.list", "localhost:9092"},
        {"group.id", group_id},
        {"enable.auto.commit", false},
    });

    auto partitions = consumer.get_partitions(topic);
    long total_lag = 0;

    for (const auto& p : partitions)
    {
        TopicPartition tp(topic, p);
        // LogEndOffset:分区最新 offset
        int64_t end_offset = consumer.query_offsets(tp).offset;
        // ConsumerOffset:该 group 已提交的 offset
        int64_t committed = consumer.committed(tp).get_offset();
        if (committed >= 0)
        {
            total_lag += end_offset - committed;
        }
    }
    consumer.close();
    return total_lag;
}

// 实际项目:把这个指标接入 Prometheus + Grafana 告警
// lag > 10000 时触发告警

7.2 死信队列(DLQ)

消息消费失败 N 次后,不再重试,转入死信队列,人工处理或告警。

RabbitMQ 死信队列配置:

#include <amqpcpp.h>

// AMQP::Channel channel(connection);

// 1. 声明死信队列(最终存放失败消息的地方)
channel.declareQueue("orders.dlq", AMQP::durable);

// 2. 声明主队列,绑定死信交换机
channel.declareQueue("orders", AMQP::durable | AMQP::Arguments{
    {"x-dead-letter-exchange",    ""},            // ★ 死信转发到默认交换机
    {"x-dead-letter-routing-key", "orders.dlq"}, // ★ 路由到死信队列
    {"x-max-retries", 3},                          // ★ 最多重试 3 次
});

Kafka 死信 Topic:

#include <sw/redis++/redis++.h>

void consume_with_dlq(Consumer& consumer, Producer& dlq_producer,
                      std::function<void(const json&)> process_func,
                      int max_retries = 3)
                      {
    Redis redis("tcp://localhost:6379");

    while (true)
    {
        Message msg = consumer.poll();
        if (!msg || msg.get_error()) continue;

        json order = json::parse(msg.get_payload());
        std::string order_id = order["order_id"];
        std::string retry_key = "retry_count:" + order_id;

        // 读取当前重试次数
        auto val = redis.get(retry_key);
        int retry_count = val ? std::stoi(*val) : 0;

        try {
            process_func(order);
            consumer.commit(msg);
            redis.del(retry_key);  // 清理重试计数
        } catch (const std::exception& e)
        {
            retry_count++;
            redis.set(retry_key, std::to_string(retry_count),
                      std::chrono::seconds(86400));

            if (retry_count >= max_retries)
            {
                // ★ 超过最大重试次数 → 发到死信 Topic
                json dlq_msg = order;
                dlq_msg["error"] = e.what();
                dlq_msg["retry_count"] = retry_count;
                dlq_producer.produce(
                    MessageBuilder("order.created.dlq").payload(dlq_msg.dump()));
                dlq_producer.flush();
                consumer.commit(msg);  // 提交了,不再重试
                std::cerr << "[死信] 订单 " << order_id << " 转入死信队列" << std::endl;
            } else
            {
                std::cerr << "[重试 " << retry_count << "/" << max_retries
                          << "] 订单 " << order_id << ": " << e.what() << std::endl;
                // 不提交,下次重新消费
            }
        }
    }
}

7.3 消费者优雅退出

收到 SIGTERM → 停止拉取新消息 → 处理完当前消息 → 提交 offset → 退出。

#include <csignal>
#include <atomic>

std::atomic<bool> running{true};

void graceful_shutdown(int)
{
    std::cout << "收到退出信号,优雅关闭中..." << std::endl;
    running = false;
}

// 注册信号处理
std::signal(SIGTERM, graceful_shutdown);
std::signal(SIGINT, graceful_shutdown);

Consumer consumer({
    {"metadata.broker.list", "localhost:9092"},
    {"group.id", "order-processor"},
    {"enable.auto.commit", false},
});
consumer.subscribe({"orders"});

while (running)
{
    Message msg = consumer.poll(std::chrono::milliseconds(100));
    if (!msg || msg.get_error()) continue;

    json order = json::parse(msg.get_payload());
    process_order(order);
    consumer.commit(msg);  // ★ 处理完才提交
}

// ★ 退出前确保最后一条消息已提交
consumer.close();

7.4 消息回溯

Kafka 可以重置 offset 回到任意位置重新消费:

// ★ 重置到最早的消息(从头消费)
Consumer consumer({
    {"metadata.broker.list", "localhost:9092"},
    {"group.id", "order-processor"},
    {"enable.auto.commit", false},
});

TopicPartition tp("orders", 0);
consumer.assign({tp});
consumer.seek(tp, 0);  // ★ offset=0 表示从头开始

while (true)
{
    Message msg = consumer.poll();
    if (!msg || msg.get_error()) continue;
    json data = json::parse(msg.get_payload());
    std::cout << "回溯消费: " << data << std::endl;
    break;  // 只收一条
}

7.5 延迟消息

RabbitMQ:死信交换机 + TTL

// 延迟队列:消息在此停留 30 秒后自动转发到消费队列
channel.declareQueue("orders.delay", AMQP::durable | AMQP::Arguments{
    {"x-dead-letter-exchange",    ""},         // ★ 死信转发到默认交换机
    {"x-dead-letter-routing-key", "orders"},   // ★ 30 秒后转发到 orders 队列
    {"x-message-ttl", 30000},                    // ★ 消息存活 30 秒
});

// 发送时也可以单独设置 TTL
json payload = {{"order_id", "A001"}};
AMQP::Message msg(payload.dump());
msg.setDeliveryMode(AMQP::DeliveryMode::persistent);
msg.setExpiration("60000");  // ★ 这条消息 60 秒后过期转发
channel.publish("", "orders.delay", msg);

Kafka:没有原生延迟,需要自己实现(时间轮 + 延迟 Topic)


8 · 更多实战场景

8.1 日志收集

各服务打日志 → Kafka → ES / Hadoop / 流计算

// 生产者:每个服务把日志发到 Kafka
Producer producer({
    {"metadata.broker.list", "localhost:9092"}
});

void log(const std::string& level, const std::string& service,
         const std::string& message)
         {
    // 替代 std::cout/std::log,日志直接进 Kafka
    json entry = {
        {"level",     level},
        {"service",   service},
        {"message",   message},
        {"timestamp", std::time(nullptr)},
    };
    producer.produce(MessageBuilder("app-logs").payload(entry.dump()));
}

8.2 事件驱动架构(EDA)

服务 A 完成操作 → 发事件 → 服务 B/C/D 响应。服务间不直接调用,通过事件协作。

// 用户注册后:发欢迎邮件 + 初始化积分 + 创建默认配置
// 传统做法:register() 里调三个服务 → 耦合
// EDA 做法:register() 只发一个事件,三个服务各自监听

std::string register_user(const std::string& username, const std::string& email)
{
    std::string user_id = create_user(username, email);
    // ★ 一个事件,N 个消费者各自响应
    json event = {
        {"user_id",  user_id},
        {"username", username},
        {"email",    email},
    };
    producer.produce(MessageBuilder("user.registered").payload(event.dump()));
    producer.flush();
    return user_id;
}

// 邮件服务监听 user.registered → 发欢迎邮件
// 积分服务监听 user.registered → 初始化积分账户
// 配置服务监听 user.registered → 创建默认配置
// ★ 新增"初始化推荐列表"功能?加个消费者就行,不用改注册代码

8.3 跨系统数据同步

数据库变更 → Canal 监听 binlog → Kafka → 下游同步到 ES / 缓存 / 数仓

图 3 \xb7 Canal 数据同步架构

// 消费 Canal 消息,同步到 Elasticsearch
// 用 cppkafka 消费 + libcurl 调 ES API
#include <cppkafka/cppkafka.h>
#include <nlohmann/json.hpp>
#include <curl/curl.h>

Consumer consumer({
    {"metadata.broker.list", "localhost:9092"},
    {"group.id", "es-sync"},
    {"enable.auto.commit", false},
});
consumer.subscribe({"mysql.changes"});

while (true)
{
    Message msg = consumer.poll();
    if (!msg || msg.get_error()) continue;

    json change = json::parse(msg.get_payload());
    std::string type = change["type"];
    std::string id = change["data"]["id"];

    if (type == "INSERT")
    {
        // POST /products/_doc/{id}
        es_index("products", id, change["data"].dump());
    } else if (type == "UPDATE")
    {
        // POST /products/_update/{id}
        es_update("products", id, change["data"].dump());
    } else if (type == "DELETE")
    {
        // DELETE /products/_doc/{id}
        es_delete("products", change["before"]["id"]);
    }
    consumer.commit(msg);
}

9 · 总结:什么时候该用,什么时候不该用

该用:

  • 需要异步化、解耦、削峰
  • 有明确的"生产者-消费者"模式
  • 可以容忍最终一致性

不该用:

  • 强一致性、实时同步场景
  • 简单的 CRUD 系统,引入 MQ 只增加复杂度
  • 团队没有运维 MQ 的能力

一句话:消息队列用好了是解耦利器,用不好是故障放大器。引入之前先想清楚:你的问题真的需要异步化吗?

9.1 学习建议

消息队列的难点不在于 API 怎么调,而在于分布式环境下的可靠性思维——消息会丢、会重复、会乱序、会积压,每个问题都需要对应的机制来兜底。建议:

  1. 先跑通最小示例——用 Docker 起 Kafka,写一个 Producer + Consumer,收发一条消息。感受一下"消息"到底长什么样
  2. 再理解可靠性——故意 kill 掉消费者,观察消息怎么办;故意发重复消息,观察幂等怎么拦截。破坏性实验比看文档有效十倍
  3. 最后看完整系统——第 6 节的订单系统把前面所有知识点串起来了:异步化、解耦、幂等、手动提交、Outbox Pattern。能看懂这个,说明入门了
  4. 实际项目中先用再调——不要一上来就追求 Exactly Once,先用 At Least Once + 幂等消费,够用了再优化

现实节奏:大多数人需要 2~3 周的实战才能真正理解消息队列。第一周跑通基本收发,第二周踩坑(重复消费、消息丢失、积压),第三周才会在项目里熟练运用。不要急。