文章

异步推送

异步推送

选型建议

场景 推荐方案 理由
进程内、事务提交后才触发 @TransactionalEventListener + @Async 事务提交后异步触发,防脏通知/脏调用,不阻塞主请求
跨进程广播、离线消息可丢 Redis Pub/Sub 无需持久化、即时送达
跨进程、轻量、有 Redis Redis Stream 不引入新组件
中小规模延迟任务、有 Redis Redis ZSet 延迟队列 不引入新组件、秒级精度够用
跨进程、高可靠、有延迟需求 RabbitMQ 延迟消息原生支持、毫秒级
超高吞吐、日志类 Kafka 吞吐王
强一致事务消息 RocketMQ 事务消息

各方案落地场景匹配

# 方案 最佳匹配场景 所在模块 核心理由(为什么是它而不是别的)
1 @TransactionalEventListener + @Async 点赞/评论/关注/收藏事务提交后 → 写 community 库 cf_notify LikeAppService / CommentAppService / FollowAppService / CollectAppService 两个价值点合一:① 防脏通知——事务回滚则通知不写(AFTER_COMMIT 兜住「事务提交前发通知」的经典坑)② @Async 独立线程池 + 拒绝策略选型 + 旁路动作——写通知派发到 notifyExecutor 线程池,不阻塞主请求与锁释放;通知是「捎个信」允许丢 → DiscardPolicy,不能丢的旁路任务 → CallerRunsPolicy 主线程兜底(拒绝策略选型详见 Q3)。
2 Redis Pub/Sub 热帖的内容部分放本地缓存,当帖子状态/内容变更(正常/隐藏/审核中、内容修改、置顶/取消置顶、逻辑删除)→ 多实例本地缓存失效广播(channel=post-cache:invalidate:{postId} PostAppService 热帖数量少放本地缓存不占内存、访问量大命中率高;状态变更必须即时生效,多实例下 A 改了 B/C 本地还旧数据 → Pub/Sub 广播清缓存。注意:只缓存内容部分(标题/正文/状态/置顶),统计数(点赞/评论/收藏/浏览)走 Redis 独立计数器 INCR,不进本地缓存、不广播。允许丢(TTL 兜底)、广播语义、零持久化 —— Pub/Sub 最原汁原味的落地场景。
3 Redis Stream 关注和发帖后投递粉丝时间线(Feed Timeline)冷启动补齐(Pull)写扩散(Push)
A 关注 B → 把 B 的历史 N 条帖子投递到 A 的 timeline
FollowAppService Feed 投递不能丢:① 大 V 近期历史帖子多,投递量大需消费组分摊 ② 投递失败要 ACK 重试,不能让用户 timeline 缺帖子 ③ 进程崩溃重启任务不能丢,需持久化
4 Redis ZSet 延迟队列 帖子置顶 72 小时自动取消置顶 PostAppService + 定时扫描服务 经典延迟任务场景:量级小(每天几千条)、秒级精度够用、项目已有 Redis。选这个场景是因为置顶是运营强感知动作,痛点好讲——人工取消容易忘,超时不取消会一直霸占置顶位。
5 RabbitMQ 帖子/评论内容审核流水线:发帖 → catfun-ai 审核 → AI 拿不准转人工审核(延迟 10min 提醒超管)→ 审核结果由 community 消费 MQ 写 cf_notify 通知作者 PostAppService / CommentAppService + catfun-ai 审核模块 审核在 catfun-ai(独立服务),结果要写 community 的 cf_notify,跨服务必须走 MQ。一个场景同时命中 RabbitMQ 的所有核心特性:autoAck=false 不丢 + TTL+DLX 延迟 + Topic Exchange 路由 + 死信兜底。审核是社区产品核心风控链路,最能体现 RabbitMQ「高可靠 + 延迟 + 路由」的综合价值。
6 Kafka 行为数据分发给分析下游:浏览 → Kafka Topic → 数据分析PV、UV各自作为独立消费组消费(落库归 community,Kafka 只做数据分发,不承担落库 community 埋点发 Kafka → 分析下游 Kafka「超高吞吐 + 日志类 + 多消费组独立消费」定位的完美映射。每秒 10 万+ 浏览、分析PV、UV各自独立消费同一数据流。注意:行为数据落库归 community,Kafka 只承担「数据分发管道」职责。
7 RocketMQ 社区积分系统:点赞(+1)/ 取消点赞(-1)→ 本地写积分流水表 + 发 MQ 通知积分系统入账,保证「DB 流水」和「积分入账」原子 LikeAppService + 积分系统 本地 DB + 服务通知的强一致原子性是 RocketMQ 事务消息独有的能力,Kafka/RabbitMQ 事务都做不到(Kafka 只管消息不管 DB)。积分/金融类场景是事务消息面试必问的典型落地。

面试考点

一、进程内异步:@TransactionalEventListener + @Async

Q1:@EventListener 和 @TransactionalEventListener 的区别?为什么本项目统一用后者?

两者用法基本一样(@Async、事件类型匹配、SpEL condition 都通用),差别只在「触发是否依赖事务」:

维度 @EventListener @TransactionalEventListener
触发时机 事件发布即触发,不管事务结果 按事务阶段触发(AFTER_COMMIT 默认)
没有事务时 ✅ 正常触发 ❌ 默认不触发(需 fallbackExecution=true
事务回滚时 ✅ 已触发(可能发脏通知) ❌ AFTER_COMMIT 不触发(防脏通知)

本项目所有「事务提交后才该做的副作用」(写 cf_notify)都依赖「事务成功才触发」,所以统一用 @TransactionalEventListener,不引入 @EventListener

1
2
3
4
5
6
// phase 可选值
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)  // 默认,最常用
//   BEFORE_COMMIT  - 事务提交前
//   AFTER_COMMIT   - 事务提交后
//   AFTER_ROLLBACK - 事务回滚后
//   AFTER_COMPLETION - 事务完成后

Q2:点赞通知为什么必须用 @TransactionalEventListener 而不是 @EventListener?

点赞操作包含「插入点赞记录 + 更新点赞数」两步,通知要写 community 库的 cf_notify。用 @EventListener 时,cf_notify 在事务提交前就写了,万一事务回滚,通知已入库 → 用户收到”点赞成功”但实际没成功。用 @TransactionalEventListener(phase = AFTER_COMMIT) 保证只有事务真正提交后才写 cf_notify


Q3:@Async 为什么不建议用默认线程池?生产环境怎么配?

默认线程池 SimpleAsyncTaskExecutor 的坑:每次调用都 NEW 一个新线程(高并发直接 OOM)、没有队列缓冲、线程名无法定位业务、拒绝策略直接抛异常。

生产级线程池配置

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
@Configuration
@EnableAsync
public class AsyncConfig {

    @Bean("notifyExecutor")
    public ThreadPoolTaskExecutor notifyExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(Runtime.getRuntime().availableProcessors() * 2);   // IO 密集型 2N
        executor.setMaxPoolSize(Runtime.getRuntime().availableProcessors() * 4);
        executor.setQueueCapacity(5000);    // 队列缓冲,防止瞬时压垮
        executor.setKeepAliveSeconds(60);   // 非核心线程空闲回收
        executor.setThreadNamePrefix("notify-async-");  // 日志可定位
        executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); // 不丢任务
        executor.setWaitForTasksToCompleteOnShutdown(true);
        executor.setAwaitTerminationSeconds(30);
        return executor;
    }
}

线程池核心线程数为什么这么配?(CPU 密集型 vs IO 密集型)

用「餐厅厨师 + 服务员」类比,N = CPU 核数:

类型 公式 类比 为什么是这个数
CPU 密集型 N + 1 厨师炒菜 4 核 → 4 个厨师刚好各占一个灶台;+1 是兜底,怕某个厨师突然被叫走(GC/页缺失),替补不上灶台空转。
IO 密集型 2N 服务员等菜 4 核 → 8 个服务员。服务员大部分时间「等后厨出菜(IO 阻塞)」,服务员多了总有一两个不用等(把 CPU 占满)。IO 阻塞越长,倍数可以越大(调外部 HTTP 接口可开到 4N~8N)。

面试一句话:CPU 密集型怕切太多,IO 密集型怕等太久。 精确公式:N * (1 + 等待时间/计算时间),2N 是经验默认值,实际靠压测调。

四种拒绝策略选型

策略 行为 适用场景
CallerRunsPolicy 调用方主线程兜底执行,不丢不抛 业务通知,不能丢的任务
AbortPolicy RejectedExecutionException 能丢的任务,系统保护型
DiscardPolicy 静默丢弃,不抛异常 完全可丢的埋点、监控采样
DiscardOldestPolicy 丢队列最老的,塞新的 “新覆盖旧”语义的场景

核心原则:不同业务用不同线程池,绝对不要所有业务共用一个大池子。


二、Redis 方案

Q4:Redis Pub/Sub 和 Redis Stream 的核心区别?

维度 Pub/Sub Stream
持久化 ❌ 不持久化,发完即丢 ✅ 持久化到内存
离线消息 ❌ 错过就丢 ✅ 可回溯历史消息
ACK ❌ 无 ACK ✅ 有消费组 + ACK
消费模式 广播(所有订阅者都收) 消费组(一条只被一个消费者处理)
适用场景 在线推送、缓存失效广播 轻量消息队列

Pub/Sub 核心缺陷:无持久化 → 断线必丢、无 ACK → 失败无法重试、慢消费者 → 缓冲区溢出被强制断开。


Q4.1:本项目用 Pub/Sub 做多实例缓存失效广播,踩过哪些坑?

本项目场景:热帖内容放 Caffeine 本地缓存,帖子被编辑/删除时,通过 Pub/Sub 广播通知所有实例清本地缓存。

坑 1:发布者只发消息不清自己的缓存 → 脏读(只能减少,无法完全避免)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
场景:3 个实例(A/B/C),负载均衡随机分发请求,热帖放 Caffeine 本地缓存

错误做法:实例 A 编辑帖子 → publish("清缓存") → 只等订阅者清
  时间线:
    ① A 更新 DB ✅
    ② A 清 Redis ✅
    ③ A publish 广播
  
  此时各实例缓存状态:
    A: [旧数据 ❌] ← 发布者自己没清!
    B: [旧数据]    ← 等广播来清,有网络延迟
    C: [旧数据]    ← 等广播来清,有网络延迟
  
  负载均衡是随机的:
    → 请求打到 A:读到旧数据 ❌(自己都没清)
    → 请求打到 B/C:广播还没到,读到旧数据 ❌

核心问题:不清自己时 A/B/C 三个实例都有脏数据(A 没清 + B/C 广播没到),清自己后 A 干净了但 B/C 在广播到达前仍有脏数据。前者脏读概率是 3/3,后者是 2/3,清自己能把脏读概率降低 1/3

正确做法:发布者发消息前先清自己的本地缓存,再 publish。

1
2
3
4
5
6
7
8
9
10
11
12
13
正确做法:实例 A 编辑帖子 → 清自己 Caffeine → publish("清缓存")
  
  此时各实例缓存状态:
    A: [已清 ✅] ← 下次读查 Redis/DB
    B: [旧数据] ← 广播还没到(网络延迟 1-5ms)
    C: [旧数据] ← 广播还没到
  
  → 请求打到 A:已清,读新数据 ✅
  → 请求打到 B/C:可能还是旧数据 ⚠️ ← 广播还没到
  
  等广播到了 B/C 之后:
    B/C: [已清 ✅]
    → 所有实例都读到新数据

注意:清自己只能减少脏读概率,无法完全消除。因为广播传播有网络延迟(哪怕 1ms),在 B/C 收到广播之前,打到 B/C 的请求还是会读到旧数据。

完全消除脏读的代价太高

方案 代价
短 TTL(如 100ms) 命中率下降,DB 压力增大
分布式读锁 读性能大幅下降
版本号比对 改表结构,缓存逻辑复杂
放弃本地缓存 失去本地缓存优势

本项目选择:「发布者清自己 + Pub/Sub 广播 + 缓存 TTL 兜底」组合策略,将脏读窗口压缩到毫秒级。帖子编辑不是高频秒级操作,短暂脏读业务可接受。金融数据场景必须用分布式锁或版本号。

坑 2:Redis 挂了 → 所有实例缓存全部失效

1
2
3
场景:Redis Pub/Sub 不可用
  实例 A 编辑帖子 → 清自己本地缓存 → publish 失败(Redis 挂了)
  → 实例 B/C 收不到消息 → 本地缓存一直是旧数据

Pub/Sub 无持久化、无重试,Redis 不可用 = 广播链路完全断。降级方案:给本地缓存设较短 TTL(如 5 分钟),即使广播丢了,TTL 到期也会重新从 DB 加载。不能只靠 Pub/Sub,TTL 兜底必须有

坑 3:慢消费者缓冲区溢出 → 被强制断开

Redis 给每个订阅者维护一个输出缓冲区,消费者处理慢了,缓冲区积压超过 client-output-buffer-limit 配置(默认 pubsub 32mb),Redis 会直接踢掉这个连接。踢掉后这个实例就收不到后续广播了,直到重连。

本项目缓存失效消息很小(就一个 postId),基本不会触发,但要知道这个限制。

面试一句话:Pub/Sub 只适合「允许丢、广播语义、消息小」的场景(如缓存失效广播),绝不能用来做业务消息传递。


Q5:Redis Stream 的消费组(Consumer Group)是什么?

消费组让多个消费者分摊同一 Stream 的消息,一条消息只被组内一个消费者处理:

1
2
3
4
5
6
Stream: msgs
  ├── Consumer Group A
  │     ├── Consumer 1 ← 处理 msg 1, 4, 7...
  │     ├── Consumer 2 ← 处理 msg 2, 5, 8...
  │     └── Consumer 3 ← 处理 msg 3, 6, 9...
  └── Consumer Group B(独立消费全部消息)

核心命令:XGROUP CREATEXREADGROUPXACKXPENDING(查未 ACK 消息)、XCLAIM(转移超时消息给其他消费者重新处理)、XTRIM(修剪历史消息)。

PEL(Pending Entries List)是消费组级别的待确认消息列表:消息投递给消费者后进入 PEL(标记属于哪个消费者),ACK 后移除。消费者宕机时消息仍在 PEL,可被其他消费者 XCLAIM 接手。


Q5.1:本项目用 Stream 做 Feed Timeline 投递,踩过哪些坑?

本项目场景:关注用户后把对方历史帖子回填到自己的时间线(Pull 模式),发帖时推送到所有粉丝时间线(Push 模式)。用 Redis Stream 异步处理,消费组多实例分摊。

坑 1:消费者名用随机 UUID → 重启后旧消息永远卡死在 PEL

1
2
3
4
5
场景:消费者名 = "fanout-consumer-" + UUID
  实例启动 → consumerName = "fanout-consumer-aaa111"
  消费了 100 条消息但没处理完就宕机 → 这 100 条留在 "aaa111" 的 PEL 里
  实例重启 → consumerName = "fanout-consumer-bbb222"(UUID 变了!)
  → "aaa111" 这个消费者名再也不会出现 → PEL 里的 100 条消息永远没人认领

这是最隐蔽的坑:消费组只认消费者名,不认进程。名字变了,旧 PEL 就是孤儿。

解决方案分两步:XCLAIM 认领消息 + XGROUP DELCONSUMER 清理死消费者

第一步:如何准确判断死消费者?

不能只看 pendingCount——活消费者刚启动还没消费时 pendingCount=0,死消费者遗留的空壳也可能 pendingCount=0,两者无法区分。

正确依据:XINFO CONSUMERS 返回的 idle 字段(距离最后一次 XREADGROUP 的时间)。

1
2
活消费者:每秒执行 XREADGROUP → idle < 2 秒
死消费者:不再读取 → idle 持续增长 → 超过 5 分钟即判定为死
idle 时间 pendingCount 判定 处理
< 5 分钟 任意 活的 跳过
>= 5 分钟 0 死的空壳 直接 XGROUP DELCONSUMER
>= 5 分钟 > 0 死的有遗留 XCLAIM 认领 → 处理 → ACK → DELCONSUMER

安全性:消费者挂了就不再读取,名下消息的 idle 也会跟着涨,一定都超阈值,XCLAIM 能全部认领。

第二步:完整清理流程

1
2
3
4
5
6
每 30 秒执行一次:
  1. XINFO CONSUMERS 查所有消费者的 idle 和 pending
  2. 跳过自己
  3. idle < 5 分钟 → 活的,跳过
  4. idle >= 5 分钟 && pending == 0 → 死的空壳,XGROUP DELCONSUMER 删除
  5. idle >= 5 分钟 && pending > 0 → XCLAIM 认领所有消息 → 处理 → ACK → XGROUP DELCONSUMER 删除

如果不清理死消费者,重启多了消费组里会堆积一大堆孤儿消费者名,占内存且让 XPENDING 扫描变慢。

坑 2:XTRIM MAXLEN 误删未消费消息 → 数据永久丢失

1
2
3
4
5
场景:设 MAXLEN 1000,消费组处理慢严重落后
  Stream 里有 1500 条消息,消费者只 ACK 到第 200 条
  执行 XTRIM MAXLEN 1000 → Redis 删掉最老的 500 条
  → 第 201~500 条消息体被删,但消费者还没消费!
  → 这些消息永久丢失,消费组再也读不到

XTRIM MAXLEN 只看条数,完全不关心是否 ACK。这是极其危险的误操作。

坑 3:消息体删了但 PEL 引用还在 → 内存不减反增

1
2
3
4
接坑 2:消息体被 XTRIM MAXLEN 删了
  但 PEL 里还保留着这些消息的元数据引用(每条至少 32 字节)
  → 消息体没了,PEL 还在
  → 内存不仅没释放,引用还是悬空的

XTRIM 只删 Stream 里的消息体,完全不碰 PEL。所以 MAXLEN 删掉未 ACK 消息 = 数据丢了 + 内存没省,双重打击。

坑 4:正确安全修剪姿势 = MINID + Pending 安全边界 + MAXLEN 兜底

1
2
3
4
5
6
7
8
9
10
错误:XTRIM MAXLEN 1000            ← 不管 ACK,乱删
错误:XTRIM MINID "1小时前的时间戳"  ← 未 ACK 的消息也被删

正确:
  1. 算时间边界:timeId = (now - 1小时) + "-0"
  2. 算 Pending 安全边界:查消费组最小 Pending ID,序列号 -1
     → 这个 ID 之前的消息才是"所有消费者都已 ACK"的安全区
  3. 取两者中更老的那个作为 XTRIM MINID 的边界
     → 有 Pending 就用 Pending 边界(保 Pending),没 Pending 就用时间边界
  4. 追加 MAXLEN=100万 兜底(极大值,只防极端内存失控,平时不触发)

核心思想:只删”所有消费组都已 ACK”的消息。Pending 最小 ID 就是安全线,线之前的才能删。

1
2
3
4
Stream: [msg1] [msg2] [msg3] [msg4] [msg5] [msg6]
                            ↑ minPendingId
        ← 安全区(可删) →  ← PEL 区(不能删)→
        XTRIM MINID 只删这里

坑 5:ACK 只是从 PEL 移除,消息体仍在 Stream → 必须定期 XTRIM

1
2
场景:消息被消费 + ACK 了,但从不 XTRIM
  → Stream 无限增长,内存爆炸

ACK 的作用只是把消息从消费者的 PEL 里移除,消息体仍然留在 Stream 里。所以即使所有消息都 ACK 了,Stream 也不会自动缩小。必须定期执行 XTRIM 清理已 ACK 的历史消息。

坑 6:Spring Data Redis 版本差异 → MINID API 不一致

版本 XTRIM MINID 支持情况
3.x opsForStream().trim(K, long) 只有 MAXLEN,无 MINID → 需 Lua 脚本发原生命令
4.x opsForStream().trim(K, XTrimOptions) 原生支持,XTrimOptions.of(TrimOptions.minId(RecordId))

本项目用 Spring Boot 4.1.0(Spring Data Redis 4.1.0),原生支持:

1
2
3
4
5
6
7
// MINID 按 ID/时间修剪
XTrimOptions minIdOptions = XTrimOptions.of(TrimOptions.minId(RecordId.of(finalMinId)));
redisTemplate.opsForStream().trim(streamKey, minIdOptions);

// MAXLEN 按条数修剪
XTrimOptions maxLenOptions = XTrimOptions.of(TrimOptions.maxLen(fallbackMaxLen));
redisTemplate.opsForStream().trim(streamKey, maxLenOptions);

面试一句话总结:Redis Stream 的可靠消费靠 PEL + XCLAIM,消费者清理靠 XINFO idle 判断 + DELCONSUMER,内存管理靠 XTRIM MINID + MAXLEN 兜底。踩过的坑:消费者名不稳定 → XCLAIM 认领 + DELCONSUMER 清理死消费者、MAXLEN 乱删 → MINID 安全边界、ACK 不删消息体 → 定期 XTRIM、版本差异 → 4.x 用 XTrimOptions。


Q6:Redis ZSet 延迟队列的实现原理和缺点?

利用 ZSet 的 score 存储消息的到期时间戳,后台线程轮询拉取到期消息:

1
2
3
4
redis.zadd("delay:queue", System.currentTimeMillis() + 5000, taskId);
// 轮询
Set<String> ready = redis.zrangeByScore("delay:queue", 0, System.currentTimeMillis());
// 原子删除(Lua 脚本)

缺点:精度依赖轮询间隔(秒级)、Redis 宕机不可用、无消费确认(需自行实现 Pending 队列)、海量任务占内存。


Q6.1:本项目用 ZSet 做帖子置顶 72 小时自动取消,踩过哪些坑?

本项目场景:管理员置顶帖子时,ZADD catfun:post:top-delay postId 到期时间戳,扫描服务每秒 ZRANGEBYSCORE 取出到期帖子自动取消置顶。

坑 1:ZRANGEBYSCORE + ZREM 非原子 → 多实例重复处理

1
2
3
4
5
场景:3 个实例同时每秒轮询 ZSet
  实例 A: ZRANGEBYSCORE 取到 [postId=123]
  实例 B: ZRANGEBYSCORE 也取到 [postId=123]   ← 同一时刻取到相同任务
  实例 A: 更新 DB isTop=0 + 清缓存
  实例 B: 也更新 DB isTop=0 + 清缓存           ← 重复处理,浪费 DB 操作

ZRANGEBYSCORE 查询和 ZREM 删除是两条独立命令,中间有时间窗口,多实例会同时取到相同任务。

解决:先 ZREM 抢占,成功才处理

1
2
3
4
5
6
7
for (String postId : expiredPostIds) {
    // ZREM 是原子的,多实例下只有一个返回 true
    if (!redisService.zRemove(POST_TOP_DELAY_KEY, postId)) {
        continue;  // 已被其他实例抢走,跳过
    }
    unpinPost(postId);  // 只有抢到的实例才处理
}

ZREM 返回值为被移除的元素数量,多实例并发时只有一个返回 1,其余返回 0。相当于用 ZREM 实现了”抢占式 ACK”。

坑 2:先 ZREM 后处理失败 → 消息丢失

1
2
3
4
场景:接坑 1 的解决方案
  ZREM 成功 → 处理失败(DB 异常、网络闪断)
  → 消息已从 ZSet 移除,但帖子还是 isTop=1
  → 永远不会再被扫描到,置顶永远不会自动取消

解决:处理失败重新入队

1
2
3
4
5
6
try {
    unpinPost(postId);
} catch (Exception e) {
    // 处理失败,重新入队 1 分钟后重试
    redisService.zAdd(POST_TOP_DELAY_KEY, postId, System.currentTimeMillis() + 60_000);
}

坑 3:无 ACK 机制 → 可靠性全靠业务自己保证

1
2
3
Redis Stream:有 PEL + XCLAIM,消费者宕机消息不丢,其他消费者能接手
Redis ZSet:  无 PEL,ZREM 就没了,宕机时正在处理的消息直接丢
RabbitMQ:    有 ACK + 重入队 + 死信队列,完整可靠性保证

ZSet 延迟队列的可靠性不如 Stream 和 MQ,只能靠”ZREM 抢占 + 失败重入队”近似实现。如果对可靠性要求高,不应选 ZSet。

坑 4:Redis 宕机 → 延迟任务全部丢失

1
2
3
4
场景:Redis 未开 AOF 或 AOF 刷盘间隔内宕机
  → ZSet 数据丢失
  → 所有未到期的置顶任务消失
  → 帖子永远置顶,不会自动取消

解决:DB 兜底 + AOF 持久化

1
2
3
-- 兜底定时任务(低频,每小时扫一次):
SELECT id FROM cf_post WHERE is_top = 1 AND top_expire_time <= NOW();
-- 扫到就执行取消置顶

AOF 持久化(appendonly yes + appendfsync everysec)最多丢 1 秒数据,配合 DB 兜底可以覆盖极端场景。

坑 5:轮询空扫 → 虽然开销小但不能忽视

1
2
3
场景:大部分时间 ZSet 里没有到期任务
  → 每秒一次 ZRANGEBYSCORE 返回空列表
  → 99% 的轮询是空扫

Redis 内存操作,单次空扫耗时微秒级,对 Redis 几乎无压力。但如果 ZSet key 数量很多(多个延迟队列),每个都每秒扫一次,累积起来也有开销。可以通过调整轮询间隔(如 5 秒)平衡精度和开销。

面试一句话总结:ZSet 延迟队列靠 score 存时间戳 + 轮询 ZRANGEBYSCORE,轻量但无 ACK。多实例并发靠 ZREM 原子抢占,处理失败靠重入队,Redis 宕机靠 DB 兜底。适合量级小、秒级精度够用、已有 Redis 的场景,不适合高可靠高吞吐场景


三、RabbitMQ:延迟消息 + Exchange + ACK

Q7:RabbitMQ 的四种 Exchange 类型和路由规则?

Exchange 类型 路由规则
Direct routing_key == binding_key(精确匹配)
Fanout 广播到所有绑定队列(忽略 routing key)
Topic routing_key 匹配通配符(* 一个词、# 零或多个词)
Headers 基于消息 headers 键值对匹配(少用)
1
2
3
4
5
6
7
Direct 示例:
  Exchange: notify-ex,绑定 notify-queue-email(key="email")、notify-queue-sms(key="sms")
  Producer 发 routing_key="email" → 只进 notify-queue-email ✅

Topic 示例:
  Exchange: order-ex,绑定 order-q-payment(key="order.paid.*")
  Producer 发 routing_key="order.paid.wechat" → 命中 order-q-payment ✅

Q8:RabbitMQ 怎么保证消息不丢?ACK 机制是什么?

生产者端三步保证

  1. Publisher Confirms:Broker 收到消息后回调 ConfirmCallback
  2. Publisher Returns:消息无法路由时回调 ReturnCallback
  3. 消息持久化:Exchange/Queue durable=true + Message deliveryMode=PERSISTENT

消费者端手动 ACK

1
2
3
4
5
6
// 处理成功 → 确认,从队列删除
channel.basicAck(deliveryTag, false);
// 处理失败 → 拒绝 + 重新入队(可能重复消费,业务要幂等)
channel.basicNack(deliveryTag, false, true);
// 处理失败 → 拒绝 + 丢弃(进死信队列)
channel.basicNack(deliveryTag, false, false);

消费者端完整保证

  1. autoAck=false(手动 ACK)
  2. 成功处理后 basicAck
  3. 异常时 basicNack(requeue 重入队,重试几次后进死信)
  4. 必须幂等消费(重入队/网络闪断都会导致重复投递)

Q9:RabbitMQ 的死信队列(DLX + DLQ)是什么?延迟消息怎么实现?

消息变成”死信”的三种情况

  1. 消费者 basicNack/basicRejectrequeue=false
  2. 消息 TTL 过期未被消费
  3. 队列达到最大长度,新消息挤成死信
1
2
3
4
5
典型架构:订单超时取消
[order-queue(TTL=30min + DLX=order-dlx-ex)]
         │ 30 分钟未支付 → TTL 到期 → 变成死信
         ▼
    [order-dlx-ex] → routing → [order-cancel-dlq] → 消费者执行取消逻辑

延迟消息两种实现方式

方式 原理 缺点
TTL + DLX 消息设 TTL,过期后进死信队列 队列级 TTL 有”头部阻塞”问题
延迟插件 插件存储到 Mnesia,到期投递 需装插件

头部阻塞:队列级 TTL 按入队顺序检查,队首 msg TTL=60s 时,后面 msg TTL=5s 也得等 60s。解决:用消息级 TTL(每条消息独立 expiration)或装延迟插件。


实战:猫趣内容审核流水线真实踩坑

以下坑全部来自猫趣 catfun-ai 内容审核流水线的真实调试经历: 发帖 → community 发 MQ → catfun-ai 消费 → Dify 审核 → 拿不准转人工(延迟 10min 提醒超管)→ 结果回传 community 写 cf_notify

Q10:消息不丢的三件套(理论见 Q8,这里对应到本项目落地)

  • 生产者端publisher-confirm(Broker 收到才回调)+ publisher-returns(路由失败回调),把消息可靠投递到交换机。
  • 服务端:Exchange / Queue / Message 都声明 durable=true(本项目队列用 QueueBuilder.durable,交换机默认 durable),Broker 重启不丢。本项目消息 deliveryMode=PERSISTENT 依赖 Spring AMQP 默认值,未显式声明,功能等价。
  • 消费端autoAck=false + 业务处理成功后再 basicAck。autoAck=true 时消息一发出就 ACK,消费者崩了消息直接丢。

Q11:autoAck=false 但消费者抛异常没 ACK 会怎样? 消息一直留在队列(Unacked 状态),RabbitMQ 不会重投给其它消费者,直到连接断开后才重新变为 Ready 重新投递。所以消费逻辑要 try-catch:成功 basicAck,失败 basicNack(requeue=true) 重投 或 requeue=false 进死信队列。切勿在循环里无限 requeue,否则消息原地打转把队列堵死。

Q12:NO_ROUTE 是什么?怎么排查? replyText=NO_ROUTE 表示交换机收到了消息,但找不到匹配 routingKey 的绑定队列。 本项目真实踩坑:延迟队列 ai.audit.manual.delay 只声明了 Queue 却漏了 Bindingcatfun.ai.direct 交换机,导致发消息到这个 routingKey 时 NO_ROUTE。 排查顺序:① 确认 BindingBuilder.bind(queue).to(exchange).with(routingKey) 是否写了;② 确认 producer 发的 routingKey 与 binding 的 routingKey 字符串完全一致;③ 确认交换机类型(direct 要求完全相等,topic 支持通配符)。

Q13:RabbitMQ 怎么实现延迟队列?(理论见 Q9,这里是本项目形态) 本项目用 TTL + DLXai.audit.manual.delay 队列是队列级 TTLx-message-ttl=10min 写在队列 args 里,所有消息统一 10 分钟过期),到期后经 DLX(catfun.ai.direct) 转投 content.audit.manual 队列提醒超管。

对比消息级 TTL:本项目是队列级,所有消息延迟相同,不存在头阻塞差异;若将来要”不同业务不同延迟”,需改用 rabbitmq-delayed-message-exchange 插件的消息级 x-delay(见 Q9 表格)。

Q14:交换机/队列/绑定必须代码声明还是手动建? 推荐两者结合:代码里用 QueueBuilder / BindingBuilder 声明(幂等,已存在则跳过),同时运维侧在管理台预建并配好权限。本项目 RabbitMqConfig 声明所有队列和绑定,启动即自动建好,避免”代码忘了建 binding 导致 NO_ROUTE”这类问题(见 Q12)。

Q15:RabbitMQ 和 Kafka 的消费者模型根本区别?

  • Kafka:消费进度由 offset 记录,消费组各自维护,支持回溯重放(重置 offset 重新消费),天然的”事件溯源”模型。
  • RabbitMQ:消息被 ACK 即从队列删除,不支持回溯,是”任务队列”模型,更适合一次性的指令/任务(如本项目的审核任务、通知)。

四、Kafka:超高吞吐 + 顺序保证

Q16:Kafka 为什么吞吐量这么高?

四大核心设计:

1. 顺序写磁盘

消息只追加到日志文件末尾,不找位置、不更新索引(RabbitMQ 维护 B+ 树索引会打断顺序写)。

1
2
随机写:磁头移到 A → 写 → 移到 B → 写 → 移到 C → 写(寻道耗时)
顺序写:磁头一直往前滑,一路写到底(无寻道)

机械硬盘上顺序写比随机写快 1000 倍,SSD 上也有数倍差距。

2. PageCache(页缓存)

Kafka 不自己管内存,直接用操作系统的空闲内存(PageCache):

  • :数据写入 PageCache 就告诉客户端”成功”,刷盘交给操作系统后台完成
  • :消费者读的消息刚好在 PageCache 里 → 直接从内存返回,不碰磁盘
  • 不占 JVM 堆:消息堆积不会触发 GC,Java 进程稳定

3. 零拷贝(Zero Copy)

1
2
3
4
5
6
传统 IO(4 次拷贝 + 4 次上下文切换):
  磁盘 → 内核缓冲区 → 用户空间(JVM) → Socket 缓冲区 → 网卡

零拷贝 sendfile(2 次拷贝 + 2 次上下文切换):
  磁盘 → 内核缓冲区 → 网卡
  (数据完全不经过 Kafka 的 JVM 内存)

好比快递从仓库直接装车发走,不需要先搬到中转站再装车。

4. 批量发送 + 压缩

linger.ms 攒一批 + lz4 压缩,减少网络请求次数和传输量。


Q17:Kafka 怎么保证消息不丢?

环节 配置 作用
Producer acks=all 等所有 ISR 副本都确认收到
Producer retries > 0 发送失败自动重试
Producer enable.idempotence=true 幂等生产者,防重试重复
Broker replication.factor >= 3 至少 3 副本,允许 1 个挂掉
Broker min.insync.replicas >= 2 最少同步副本数
Broker unclean.leader.election.enable=false 禁止非 ISR 副本当选 Leader
Consumer enable.auto.commit=false 手动提交 offset,处理完再提交

Q18:Kafka 分区能保证顺序性吗?破坏顺序的核心陷阱有哪些?

结论:Kafka 只保证「同一分区内」有序,分区间无顺序保证。

三大破坏顺序的高频陷阱

# 陷阱 原理 解决方案
1 多分区乱序 不同分区之间无全局时钟,写入和消费都是并行的 按业务 key 哈希路由到同一分区
2 生产者重试乱序 retries>0max.in.flight>1 时,先发的请求超时重试、后发的请求先成功,offset 颠倒 开启幂等生产者 enable.idempotence=true,或 in.flight=1
3 消费者异步/线程池乱序 poll 后直接丢线程池并发处理,同 key 的消息被不同线程乱序执行 消费端按 key 分桶排队,同一 key 交给同一业务线程串行处理

生产者重试乱序原理

1
2
3
4
5
6
7
Producer 发同一 key 的 3 条消息(in-flight=5):
  [batch1(创建)] → 发出 ✅
  [batch2(支付)] → 发出 ✅
  [batch3(发货)] → 发出 ✅

Broker 响应:batch3 先成功 → batch2 成功 → batch1 超时重试后才成功
→ Broker 上的 offset 实际顺序变成:发货(0) → 支付(1) → 创建(2) ❌

幂等生产者解决原理: Producer 申请 PID,每条消息带 (PID, Partition, Seq Number),Broker 维护最大 Seq,Seq 不递增的重试消息直接忽略。

生产者配置

1
2
3
enable.idempotence = true
max.in.flight.requests.per.connection = 5
acks = all

消费者正确姿势

  • enable.auto.commit=false先处理,后 commit
  • 消费端按 key 分桶(hash(key) % N),每个桶用单线程池串行执行
  • 处理失败不 commit,转死信兜底

面试一句话总结:Kafka 的顺序性是 Producer + Broker + Consumer 三层协同的结果,只说”用 key 路由”只答对了 1/3。


Q19:Kafka 的 ISR、HW、LEO 是什么?

用「老师抄笔记」做比喻:

1
2
3
4
5
场景:老师在黑板上写笔记,2 个学生同步抄写(1 老师 + 2 学生)

「老师」= Leader 主节点(写消息的人)
「学生」= Follower 从节点(同步数据的人)
「笔记第 N 行」= offset(消息偏移量)

ISR(靠谱学生名单)

1
2
3
4
不是所有学生都能跟上。有些学生上课玩手机,抄的比老师写的慢 30 秒以上,
就会被踢出「靠谱学生名单」——只有名单里的学生才认可是合格的。

ISR = In-Sync Replicas = 和老师同步的靠谱学生集合

LEO(老师写到第几行了)

1
2
3
4
5
6
7
8
9
10
老师在黑板上已经写了 8 行(第 0-7 行)
→ 老师的 LEO = 7

学生 A 抄到了第 6 行(抄了 0-5,正在写第 6 行)
→ 学生 A 的 LEO = 6

学生 B 也抄到了第 6 行
→ 学生 B 的 LEO = 6

LEO = Log End Offset = 我自己已经抄到/写到了第几行

HW(全班都抄到的最低行)

1
2
3
4
5
6
7
8
9
HW = High Watermark = 全班都抄到了的那个行号

    老师(已写):0 1 2 3 4 5 6 7   → LEO = 7
    学生 A(已抄):0 1 2 3 4 5       → LEO = 6
    学生 B(已抄):0 1 2 3 4 5       → LEO = 6
                              ↑
                    HW = 6(取 ISR 里最小的 LEO)

也就是说:前 6 行(0-5)全班都抄完了,第 6-7 行只有老师写了,学生还没抄完

消费者只读 HW 之前的内容

1
2
3
4
5
6
班长(消费者)想抄一份笔记拿去看:
  看第 0-5 行 ✅ → 全班都抄完了,就算老师请假了,学生那里也有备份
  看第 6-7 行 ❌ → 只有老师黑板上有,万一老师的黑板擦了(主节点挂了)
                    学生上位当老师,那两行就没了 → 班长看到的就是"假数据"

所以 Kafka 规定:消费者绝对不能读 HW 之后那些「只老师有、学生还没抄完」的脆弱数据。

一句话总结:

名词 大白话 比喻
ISR 靠谱学生名单 跟老师同步速度不超过 30 秒的学生
LEO 我写到第几行了 老师/学生各自的进度
HW 全班都抄完的最低行 取所有靠谱学生进度中最小的那个

Broker 配置三件套(防丢数据)

1
2
3
replication.factor=3          → 至少 3 份备份(1 老师 + 2 学生)
min.insync.replicas=2         → acks=all 时,至少 2 个人抄完才算写成功
unclean.leader.election=false → 上课玩手机被踢出名单的学生,绝不能当临时老师

Q20:Kafka Rebalance 是什么?有什么危害?怎么减少影响?

Rebalance 是 Consumer Group 内消费者实例变化时,重新分配 partition 的过程。

触发条件:新增消费者、消费者宕机、业务处理超时(超过 max.poll.interval.ms)、Topic 分区扩容。

危害

  1. STW(Stop-The-World):Rebalance 期间整个 Consumer Group 停止消费
  2. 重复消费:未 commit offset 的消息重放
  3. 消费延迟飙升:积压消息集中涌入

减少影响

  1. max.poll.interval.ms 设得比业务最大处理时间大
  2. group.instance.id 固定消费者身份,重启不触发 Rebalance
  3. 异步解耦:消费者只拉消息,业务逻辑丢给独立线程池,避免 poll 被阻塞

Q21:Kafka 幂等生产者和事务生产者的区别?

用「寄快递」做比喻:


幂等生产者:同一个包裹寄两次,收件人只收到一个

1
2
3
4
5
6
7
8
场景:你寄了一个包裹,网络不好没收到回执,你以为没寄出去又寄了一次。

没有幂等:收件人收到 2 个一模一样的包裹 ❌
有幂等:  快递员发现"这不是刚寄过的吗",第二次直接拒收 ✅

原理:Kafka 给每个生产者发一个 PID(快递员工号),
      每条消息带一个递增的 SeqNumber(包裹编号)。
      Broker 收到后发现"这个编号我收过了",直接丢弃重复的。

局限:只在「同一个分区」内有效。如果消息被路由到不同分区,编号就对不上了。


事务生产者:要么全寄出去,要么一个都不寄

1
2
3
4
5
6
7
8
9
场景:你要给 3 个人分别寄合同、发票、收据,必须同时寄出,不能只寄了合同没寄发票。

没有事务:合同寄了,寄发票时网络断了 → 收件人只收到合同 ❌
有事务:  3 个包裹一起打包,全部贴好才一起寄出,任何一个没贴好就全部撤回 ✅

原理:Kafka 用两阶段提交(2PC):
      1. begin → 写消息到多个分区(此时消费者看不到)
      2. commit → 所有分区同时标记为"可消费"(消费者才能看到)
         abort  → 所有分区回滚(消费者永远看不到)

对比

维度 幂等生产者 事务生产者
解决什么问题 重试导致重复 跨分区不原子
比喻 同一包裹寄两次,只收一个 多个包裹要么全寄要么全不寄
范围 单分区 多分区
配置 enable.idempotence=true transactional.id=xxx
本项目用哪个 ✅ 用了幂等(防重试重复) ❌ 没用(不需要跨分区原子)

重要:Kafka 事务只管「消息之间原子」,不管 DB。比如”扣库存 + 发消息”这种”DB 操作 + MQ 操作”要原子的场景,Kafka 事务帮不了你,得用 RocketMQ 事务消息。


Q21.1:本项目用 Kafka 做浏览数据 PV/UV 统计,踩过哪些坑?

本项目场景:用户浏览帖子 → community 发 Kafka 浏览事件 → pv-group 和 uv-group 两个消费组各自消费 → 更新 Redis ZSet 排行榜 → 产出 PV TOP10 / UV TOP10 接口。

坑 1:消费者重平衡导致重复消费 → PV 多算

1
2
3
4
场景:uv-group 消费者实例重启,触发 Rebalance
  → 分区重新分配,未 ACK 的消息被重新投递
  → 同一条浏览消息被消费两次
  → PV ZINCRBY +1 两次 → PV 多算 1

解决:eventId 幂等去重

1
2
3
4
5
6
7
// 每条消息带唯一 eventId(雪花 ID)
// 消费前 SADD 去重 Set,返回 false 说明已处理过
if (!redisService.sAdd("post:pv:dedup", eventId)) {
    ack.acknowledge();  // 已处理,直接 ACK 跳过
    return;
}
redisService.zIncrementScore("post:pv:rank", postId, 1.0);

坑 2:分区策略选错 → 同一帖子 PV 计数乱序

1
2
3
4
场景:分区策略用 RoundRobin(轮询)
  → 同一帖子 postId=123 的两条浏览消息被分到不同分区
  → 两个分区被不同消费者线程并发消费
  → PV 计数顺序不可控(虽然 ZINCRBY 是原子的,但消费顺序乱了影响调试和排查)

解决:按 postId Hash 分区

1
2
3
// 生产者发送时 postId 作为 key
kafkaMessageProducer.sendAsync(TOPIC_VIEW, postId, JSONUtil.toJsonStr(event));
// Kafka 生产者默认按 key Hash 分区,同一 postId 路由到同一分区

坑 3:消费处理慢 → 超过 max.poll.interval.ms 被踢出消费组

1
2
3
4
场景:Redis 不可用,消费处理卡住
  → 超过 max.poll.interval.ms(默认 30s)
  → 被踢出消费组,触发 Rebalance
  → 其他消费者接手,但正在处理的消息未 ACK → 重复消费

解决:死信兜底 + 合理的超时配置

1
2
3
4
5
6
7
8
9
10
try {
    // 业务处理
    redisService.zIncrementScore(...);
    ack.acknowledge();
} catch (Exception e) {
    // 处理失败,ACK 跳过(不阻塞正常消费)
    // 失败消息可通过日志补捞,或发到死信 Topic
    log.error("PV 消费失败 eventId={}", eventId);
    ack.acknowledge();
}

配置上 max.poll.interval.ms=30000,确保单条消息处理不会超过 30 秒。

坑 4:两个消费组消费同一份消息 → 各自维护 offset,互不影响

1
2
3
4
pv-group:消费消息 → ZINCRBY post:pv:rank(每次 +1)
uv-group:消费消息 → SADD post:uv:{postId} {userId} → 新用户才 ZINCRBY post:uv:rank

两个消费组各自维护消费进度(offset),一个组落后不影响另一个组。

这是 Kafka 多消费组的天然能力,不需要额外处理。

坑 5:NewTopic Bean 声明了 3 分区,监控却显示 1 分区

1
2
3
4
5
6
7
场景:BehaviorEventProducer 中定义了 NewTopic Bean(分区数=3)
  @Bean
  public NewTopic viewTopic() {
      return new NewTopic("community-behavior-view", 3, (short) 1);
  }

启动服务后 Kafka 监控显示分区数仍然是 1

原因:Spring 的 KafkaAdmin 扫描到 NewTopic Bean 后,只负责创建不存在的 Topic。如果 Topic 已经存在(之前以默认 1 分区被创建过),KafkaAdmin 会直接跳过,不会修改已存在 Topic 的分区数

解决:启动时用 AdminClient 自动检查并扩容

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
@Component
public class KafkaTopicPartitionInitializer implements CommandLineRunner {

    private final Map<String, NewTopic> newTopicMap;
    private final AdminClient adminClient;

    @Override
    public void run(String... args) {
        // 1. 查询 Kafka 中已存在的所有 Topic
        Set<String> existing = adminClient.listTopics().names().get(10, TimeUnit.SECONDS);

        for (NewTopic topic : newTopicMap.values()) {
            if (!existing.contains(topic.name())) {
                // 不存在 → 创建
                adminClient.createTopics(List.of(topic)).all().get();
            } else {
                // 已存在 → 查实际分区数,不够则扩容
                int actual = describePartitions(topic.name());
                if (actual < topic.numPartitions()) {
                    adminClient.createPartitions(Map.of(
                        topic.name(), NewPartitions.increaseTo(topic.numPartitions())
                    )).all().get();
                }
            }
        }
    }
}

坑 6:describeTopics 对不存在的 Topic 直接抛异常

1
2
3
场景:先用 describeTopics 查所有 Topic 状态,再判断是否为 null
  → Topic 不存在时 describeTopics 直接抛 UnknownTopicOrPartitionException
  → 不是返回 null,整个检查流程中断

原因describeTopics 要求传入的 Topic 名称必须已存在,不存在直接抛异常,不像 listTopics 那样能优雅处理。

解决:先用 listTopics().names() 拿到所有已存在 Topic 的名称集合,判断存在性后再决定是创建还是扩容。

注意事项

  • Kafka 只支持增加分区,不支持减少分区,减少分区只能删除 Topic 重建
  • 扩容后已有消息不会重新分区,旧消息仍在原分区,只有新消息按新分区数路由
  • 扩容失败不应阻塞服务启动,用 try-catch 兜住

面试一句话总结:Kafka 消费可靠性靠手动 ACK + 幂等去重,多消费组靠独立 offset 互不影响,分区靠 key Hash 保序,重平衡靠 CooperativeStickyAssignor 减少影响,失败靠死信兜底,Topic 分区扩容靠 AdminClient 启动时自动检查


五、RocketMQ:事务消息 + 顺序消息 + 延时消息 + Outbox

基于社区积分系统(catfun-community)实践,覆盖事务消息、顺序消息、延时消息、普通消息四种用法。

Q22:项目中 RocketMQ 的四种消息用法?

用法 场景 代码位置
事务消息 积分流水写入与消息投递原子化 PointTransactionProducer
顺序消息(消费端) 同一 Queue 内串行消费积分变动 PointChangeConsumerConsumeMode.ORDERLY
延时消息 消费失败后延迟重试入账 PointChangeConsumer.sendRetrysyncSendDelayTimeSeconds
普通并发消息 重试 Topic 的最终兜底消费 PointRetryConsumer

Q23:RocketMQ 事务消息的流程?(两阶段 + 回查)

两阶段提交:

1
2
3
4
5
6
7
① Producer 发送 Half Message(Broker 存储但不投递)
② Broker 返回成功 → 回调 executeLocalTransaction()
③ 本地事务执行(写 cf_point_flow,标记 COMMITTED)
④ 返回 COMMIT / ROLLBACK / UNKNOWN
    ├── COMMIT   → Broker 投递消息到原始 Topic
    ├── ROLLBACK → Broker 丢弃半消息
    └── UNKNOWN  → Broker 定时回查 checkLocalTransaction()

回查机制:

Broker 默认每 60 秒回查一次,最多回查 15 次。回查时查 cf_point_flow.tx_status

tx_status 返回值 含义
COMMITTED COMMIT 本地事务已成功,投递消息
ROLLBACK ROLLBACK 本地事务失败,丢弃消息
PENDING / 不存在 UNKNOWN 还没执行完,等下一轮

项目中的实现: PointTransactionProducer.java

  • sendPointChange():构建 payload,调 rocketMQTemplate.sendMessageInTransaction(topic, msg, userId)
  • executeLocalTransaction():解析 payload → 调 PointAppService.commitFlowAsCommitted() 写流水
  • checkLocalTransaction():调 PointAppService.checkFlowTxStatus(flowId) 查 tx_status

为什么不用 Kafka 事务? Kafka 事务是生产者级事务,只管「多 Topic 消息原子性」,完全不涉及业务本地事务(扣库存、写订单表这些 DB 操作)。RocketMQ 事务消息才能保证「本地 DB 操作 + MQ 消息可见」原子。


Q24:事务消息能保证顺序消费吗?

结论:不能。事务消息和顺序消息在 API 层面互斥。

原因: 事务消息内部调用链走的是 send(msg)MQFaultStrategy.selectOneMessageQueue() 轮询选队列,不经过 MessageQueueSelector;而顺序消息需要 syncSendOrderlysend(msg, MessageQueueSelector, hashKey) 按 hashKey 路由。两条 API 路径完全独立。

当前 sendMessageInTransaction(topic, msg, userId)userId 参数只透传给 executeLocalTransaction 作为回调参数,不参与队列路由。同一用户的消息会被轮询分散到不同 Queue。

怎么兜底? 积分正确性不依赖消费顺序,靠原子 SQL + 幂等去重

  • 原子 SQLUPDATE cf_user_point SET total_point = total_point + ? WHERE user_id = ?,数据库引擎通过行锁保证并发加减的正确性,多条消息同时到达同一用户也不会算错
  • flowId 幂等去重:每条消息携带唯一 flowId(雪花 ID),消费前检查是否已处理过,防止重复消费导致积分多加/多减
  • level 覆盖容忍:level 是 total_point 的衍生值,并发下先查再算再写可能短暂不准,但不影响积分本身,下一次消费基于最新 total 会修正。举例:
    1
    2
    3
    4
    5
    6
    7
    
    A: increment +1 提交 → total=100
    A: findByUserId → 100 → calcLevel → Lv1
    B: increment +1 提交 → total=101
    B: findByUserId → 101 → calcLevel → Lv2
    A: save level=1
    B: save level=2  或 A 后写,结果可能 level=1(被旧值覆盖)
    → 下次消费时读到 total=101,重新算出 Lv2 并写回,完成修正
    

彻底的兜底方案:一条 UPDATE 同时改 total + level。 把等级分段用 CASE WHEN 嵌进原子 SQL,一条语句同时更新 total_point 和 level,从根本上消除先查再写的并发问题:

1
2
3
4
5
6
7
8
9
10
11
UPDATE cf_user_point
SET total_point = total_point + #{delta},
    level = CASE
        WHEN (total_point + #{delta}) <= 100 THEN 1
        WHEN (total_point + #{delta}) <= 500 THEN 2
        WHEN (total_point + #{delta}) <= 2000 THEN 3
        WHEN (total_point + #{delta}) <= 10000 THEN 4
        ELSE 5
    END,
    update_time = NOW()
WHERE user_id = #{userId}

Q25:RocketMQ 怎么保证顺序消息?

前提:这是普通顺序消息的方案,不是事务消息

两层配合,缺一不可:

机制 API 作用
发送端 MessageQueueSelector 按 hashKey 路由 syncSendOrderly(topic, msg, hashKey) 同一 hashKey → 同一 Queue
消费端 MessageListenerOrderly 串行消费 consumeMode = ConsumeMode.ORDERLY 同一 Queue 内单线程消费
  • 发送端不路由 → 同一用户散落到多 Queue → 多线程并行 → 乱序
  • 消费端不 ORDERLY → 同一 Queue 多线程并发拉取 → 乱序

消费端是”接力者”不是”决策者”:ConsumeMode.ORDERLY 只保证同一 MessageQueue 内串行,跨 Queue 的顺序完全不管。真正的路由决策在发送端。

如果项目需要事务 + 顺序怎么办? 改用 Outbox 模式

1
2
3
4
① 本地事务:写 cf_point_flow(PENDING)
② 事务提交后:syncSendOrderly(topic, msg, userId)  ← hashKey 路由
③ 消费端:ORDERLY 顺序消费,入账后标记 flow 为 CONSUMED
④ 定时补偿:扫描超时仍 PENDING 的流水,重新投递

事务靠本地事务 + 定时补偿兜底,顺序靠 syncSendOrderly + ORDERLY 两层保证。

如果不依赖延时消息且非金融级安全业务,建议换 Kafka。 Outbox 模式本身已经绕开了 RocketMQ 最核心的事务消息特性,剩下的需求是”顺序 + 最终一致 + 补偿”。这种场景 Kafka 的语义更干净,性能上也更有优势

维度 RocketMQ(Outbox) Kafka(Outbox)
顺序语义 syncSendOrderly(hashKey) 选 Queue + 消费端 ORDERLY 串行,手动实现两层 Partition 是天然顺序单元,按 key hash 到同 Partition,同一 Partition 只分配给组内一个线程,模型原生支持
生产吞吐 syncSendOrderly 是同步发送,每条要等 Broker ACK 返回;同一 Queue 还需等前一条落盘确认,QPS 上限低 客户端批量 + 压缩 + 异步刷盘,生产端吞吐通常是 RocketMQ 的 2-5 倍,顺序消息场景差距更明显
消费吞吐 ORDERLY 模式下同一 Queue 只有单线程消费,Queue 数受限(通常 4-16),消费并发扩展靠加 Queue Partition 数可以轻松开到 64-200,同一消费组内按 Partition 分配线程,消费并发度线性扩展,消费吞吐差距显著
存储模型 CommitLog + ConsumeQueue 两层存储,随机读放大;顺序消费时还要按 Queue 索引定位 仅用 Partition 日志文件顺序追加,零拷贝 + page cache,顺序读写下 IO 路径更短
延时消息 可以用 RocketMQ 5.x 任意秒级延时,补偿时间控制更精准 没有 Broker 级延时消息,需在流水表加 next_retry_at,由定时扫描决定重发时机(Outbox 的定时扫描本来就要扫,不额外增加成本

Outbox 模式下安全维度对比:

维度 RocketMQ(严格配置) Kafka(严格配置)
消息零丢失 flushDiskType=SYNC_FLUSH + brokerRole=SYNC_MASTER(同步刷盘 + 同步复制) acks=all + min.insync.replicas=2 + replication.factor=3 + unclean.leader.election.enable=false
顺序保障链路 syncSendOrderly + ORDERLY 两层手动实现,链路长出错概率高 key → Partition → Consumer 线程映射,模型原生,链路短出错概率低
死信/消费失败兜底 内置 %DLQ%topic 死信队列,超限自动进 需业务自己实现(发独立 DLQ Topic / 写 DB 失败表)
消息审计/追溯 内置 MessageTrace 轨迹开关,mqadmin 直接查消息轨迹 需额外链路追踪组件(OpenTelemetry 等),没有 Broker 原生轨迹
误删/误发恢复 支持重置消费位点(resetOffsetByTime),可按时间回溯 支持重置消费位点(--to-datetime),可按时间回溯

结论: 如果从零选型且场景是 Outbox + 严格顺序 + 高吞吐要求 + 非金融级安全业务,优先选 Kafka;如果是金融级业务(资金余额、支付扣款等),国内优先选 RocketMQ(运维生态成熟、案例多、有 Broker 级事务消息可不用 Outbox),海外选 Kafka。


Q26:RocketMQ 延时消息怎么用?

4.x:固定 18 个延迟等级

1
1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h

API:syncSend(topic, msg, timeout, delayLevel)

5.x:任意延时(rocketmq-spring 2.3.6 支持)

API:syncSendDelayTimeSeconds(topic, msg, delaySeconds)

项目中的使用: 消费失败后发延时消息到 RETRY Topic,10 秒后重试。


Q27:消费失败怎么处理?

两级重试机制:

1
2
3
4
5
6
7
8
9
10
11
PointChangeConsumer(顺序消费)
  ├── 消费成功 → 正常返回(自动 ACK)
  └── 消费失败 → sendRetry() 发延时消息到 RETRY Topic → 返回成功(不阻塞队列)
        │
        ▼
PointRetryConsumer(普通并发消费 RETRY Topic)
  ├── 重试成功 → 正常结束
  └── 重试失败 → 抛异常 → RocketMQ 自动 RECONSUME_LATER
        │
        ▼(maxReconsumeTimes = 5 次后)
%DLQ%topic 死信队列(人工兜底处理)

为什么顺序消费者不直接抛异常重试? 顺序消费抛异常会触发 RocketMQ 内部重试,阻塞当前 Queue 后续所有消息。改为发延时消息到独立 RETRY Topic,当前 Queue 可以继续消费后续消息,不阻塞其他积分变动。


Q28:RocketMQ 事务消息和本地事务表(Outbox)怎么选?

对比项 事务消息 Outbox 模式
原子性保证 Broker 半消息 + 回查 本地事务 + 定时补偿
顺序消费 ❌ 不支持(轮询路由) ✅ 支持(syncSendOrderly)
实现复杂度 低(框架封装好) 中(需定时任务 + 补偿逻辑)
依赖 强依赖 Broker 可用 弱依赖(本地事务先成功)
适用场景 不需要顺序的分布式事务 需要顺序 + 最终一致

项目选择: 当前用事务消息,因为积分正确性靠原子 SQL 保证,不依赖消费顺序。如果未来有需要严格顺序的场景(如账户余额增减),再切换为 Outbox 模式。


Q29:消息队列(Queue)的数量由谁决定?

Topic 创建时的运维配置决定,生产端和消费端代码都不能决定。

层级 配置方式 说明
Broker 默认值 broker.confdefaultTopicQueueNums 自动创建 Topic 时用此默认值(通常 4 或 8)
显式创建 mqadmin updateTopic -w 8 -r 8 -w 写队列数、-r 读队列数
动态修改 mqadmin updateTopic 再次执行 可在线增加,不会减少

生产者只能从已有队列列表中选一个发送,消费者按已有队列数分配消费线程。


Q30:RocketMQ 事务消息的循环依赖怎么解决?

1
2
3
PointAppService → (依赖) PointMessagePort
PointTransactionProducer → (实现) PointMessagePort
PointTransactionProducer → (回查需要) PointAppService

构造器注入会形成循环依赖。解决: PointTransactionProducerPointAppService 改用 setter 注入,Spring 创建 Bean 时先用构造器实例化,再通过 setter 注入打破循环。


Q31:幂等性怎么保证?

基于 flowId 去重:每条积分变动消息携带唯一的 flowId(雪花 ID),消费前检查 cf_point_flow 是否已消费。消息重复投递或消费成功但 ACK 丢失时,flowId 已存在则直接跳过。

注意:UPDATE total_point = total_point + ? 本身不是幂等的,重复执行会导致积分多加,必须配合 flowId 去重使用。


六、综合选型

Q32:点赞发通知,用 MQ 还是进程内事件?

通知不独立服务、cf_notify 归 community 库。按「通知来源是否跨服务」选:

场景 推荐方案
community 内行为(点赞/评论/关注/收藏)→ 通知同在 community @TransactionalEventListener 直接写 cf_notify,无需 MQ
跨服务行为(catfun-ai 审核结果、积分系统入账)→ 要写 community 的 cf_notify RabbitMQ 投递,community 消费写库
超高并发(百万级点赞) Kafka 削峰 + 异步消费写 cf_notify
事务要求强一致(本地 DB + MQ 原子) RocketMQ 事务消息

关键原则:通知跟着业务领域走,不抽独立通知服务。community 内的通知直接写本库 cf_notify;跨服务的通知才走 MQ,且消费方仍是 community 自己写 cf_notify

Q33:如果 community 服务挂了,各方案的消息会丢吗?

方案 挂了会丢吗 原因
@TransactionalEventListener ✅ 丢 进程内、无持久化
Redis Pub/Sub ✅ 丢 无持久化
Redis Stream ❌ 不丢 持久化 + 消费组
RabbitMQ ❌ 不丢 持久化 + ACK
Kafka ❌ 不丢 持久化 + ISR
RocketMQ ❌ 不丢 持久化 + 同步刷盘可选

Q34:秒杀场景该用什么?

Kafka 削峰 + Redis 预扣库存 + 数据库最终落库

  1. 请求先进 Kafka(削峰,避免直接打垮 DB)
  2. 消费者从 Kafka 读消息,Redis 预扣库存(DECR 原子操作)
  3. 异步落库到 MySQL(批量 INSERT)
  4. 超时未支付的库存回补(延迟消息或定时任务)

Q35:本项目所用异步技术对比与用处总结

1. 技术清单与定位

技术 类型 在猫趣的用途 核心价值
@TransactionalEventListener + @Async 进程内异步 点赞/评论/关注/收藏事务提交后写 cf_notify 防脏通知 + 不阻塞主请求
Redis Pub/Sub 跨进程广播 多实例本地缓存失效广播 即时、可丢、零持久化
Redis Stream 跨进程持久化队列 关注/发帖后写扩散 Feed Timeline 不丢 + ACK 重试 + 消费组
Redis ZSet 延迟队列 延迟任务 帖子置顶 72h 自动取消 不引新组件、秒级精度
RabbitMQ 可靠消息队列 内容审核流水线(发帖→AI审核→转人工延迟提醒→结果回写) 高可靠 + TTL/DLX 延迟 + 路由 + 死信兜底
Kafka 高吞吐日志管道 行为数据分发(PV/UV 各自消费组) 超高吞吐 + 多消费组独立消费 + 可回溯
RocketMQ 事务消息 积分系统(DB 流水 + 积分入账原子) 事务消息保证本地 DB 与消息原子

2. 选型决策树(一句话记忆)

  • 要不要跨进程? 否 → 进程内事件;是 → 看下面。
  • 能不能丢? 能丢 + 即时 → Redis Pub/Sub;不能丢 → RabbitMQ/Kafka/RocketMQ/Redis Stream。
  • 要不要延迟? 要 + 有 Redis → ZSet;要 + 高可靠 → RabbitMQ TTL+DLX。
  • 吞吐大不大? 日志/行为类超高吞吐 → Kafka;强一致事务 → RocketMQ。
  • 跨服务写本库通知? → RabbitMQ 投递、消费方自己写库。

3. 为什么内容审核用 RabbitMQ 而不是别的

  • 不是 Kafka:审核是任务型(每条处理完即从队列删除),不是日志流,不需要回溯重放;Kafka 的 offset 模型在这里是浪费。
  • 不是 Redis Stream:审核链路要 autoAck=false 的可靠确认 + 死信兜底,Redis Stream 的 ACK 语义偏轻,且猫趣已有 RabbitMQ。
  • 不是 RocketMQ:审核不需要”本地 DB + 消息”的强一致原子(审核结果落库由消费方独立事务保证幂等即可),用不到事务消息。
  • RabbitMQ 正好集齐审核要的全部特性:持久化不丢、TTL+DLX 延迟转人工、direct 路由分流(request/result/manual/delay)、死信队列兜底、自动重连自愈。一个场景把 RabbitMQ 的看家本领全用上了。

4. 一句话总结各技术”最不可替代”的点

  • 进程内事件:事务提交后才触发(防脏数据),别的都做不到。
  • Redis Pub/Sub:最轻的跨进程广播,零持久化、零依赖(有 Redis 就行)。
  • Redis Stream:轻量但不丢的队列,ACK + 消费组,不引新组件。
  • Redis ZSet:最省事的延迟任务,已有 Redis 时首选。
  • RabbitMQ:可靠任务队列 + 延迟 + 路由三合一,企业级风控链路首选。
  • Kafka:吞吐 + 可回溯,数据分发/日志流霸主。
  • RocketMQ:事务消息唯一能绑本地 DB 的,金融/积分场景必选。

本文由作者按照 CC BY 4.0 进行授权