异步推送
异步推送
选型建议
| 场景 | 推荐方案 | 理由 |
|---|---|---|
| 进程内、事务提交后才触发 | @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 CREATE、XREADGROUP、XACK、XPENDING(查未 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 机制是什么?
生产者端三步保证:
Publisher Confirms:Broker 收到消息后回调ConfirmCallbackPublisher Returns:消息无法路由时回调ReturnCallback- 消息持久化: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);
消费者端完整保证:
autoAck=false(手动 ACK)- 成功处理后
basicAck - 异常时
basicNack(requeue 重入队,重试几次后进死信) - 必须幂等消费(重入队/网络闪断都会导致重复投递)
Q9:RabbitMQ 的死信队列(DLX + DLQ)是什么?延迟消息怎么实现?
消息变成”死信”的三种情况:
- 消费者
basicNack/basicReject且requeue=false - 消息 TTL 过期未被消费
- 队列达到最大长度,新消息挤成死信
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 却漏了 Binding 到 catfun.ai.direct 交换机,导致发消息到这个 routingKey 时 NO_ROUTE。
排查顺序:① 确认 BindingBuilder.bind(queue).to(exchange).with(routingKey) 是否写了;② 确认 producer 发的 routingKey 与 binding 的 routingKey 字符串完全一致;③ 确认交换机类型(direct 要求完全相等,topic 支持通配符)。
Q13:RabbitMQ 怎么实现延迟队列?(理论见 Q9,这里是本项目形态)
本项目用 TTL + DLX,ai.audit.manual.delay 队列是队列级 TTL(x-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>0 且 max.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 分区扩容。
危害:
- STW(Stop-The-World):Rebalance 期间整个 Consumer Group 停止消费
- 重复消费:未 commit offset 的消息重放
- 消费延迟飙升:积压消息集中涌入
减少影响:
max.poll.interval.ms设得比业务最大处理时间大group.instance.id固定消费者身份,重启不触发 Rebalance- 异步解耦:消费者只拉消息,业务逻辑丢给独立线程池,避免 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 内串行消费积分变动 | PointChangeConsumer(ConsumeMode.ORDERLY) |
| 延时消息 | 消费失败后延迟重试入账 | PointChangeConsumer.sendRetry → syncSendDelayTimeSeconds |
| 普通并发消息 | 重试 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;而顺序消息需要 syncSendOrderly → send(msg, MessageQueueSelector, hashKey) 按 hashKey 路由。两条 API 路径完全独立。
当前 sendMessageInTransaction(topic, msg, userId) 的 userId 参数只透传给 executeLocalTransaction 作为回调参数,不参与队列路由。同一用户的消息会被轮询分散到不同 Queue。
怎么兜底? 积分正确性不依赖消费顺序,靠原子 SQL + 幂等去重:
- 原子 SQL:
UPDATE 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.conf 的 defaultTopicQueueNums |
自动创建 Topic 时用此默认值(通常 4 或 8) |
| 显式创建 | mqadmin updateTopic -w 8 -r 8 |
-w 写队列数、-r 读队列数 |
| 动态修改 | mqadmin updateTopic 再次执行 |
可在线增加,不会减少 |
生产者只能从已有队列列表中选一个发送,消费者按已有队列数分配消费线程。
Q30:RocketMQ 事务消息的循环依赖怎么解决?
1
2
3
PointAppService → (依赖) PointMessagePort
PointTransactionProducer → (实现) PointMessagePort
PointTransactionProducer → (回查需要) PointAppService
构造器注入会形成循环依赖。解决: PointTransactionProducer 对 PointAppService 改用 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 预扣库存 + 数据库最终落库:
- 请求先进 Kafka(削峰,避免直接打垮 DB)
- 消费者从 Kafka 读消息,Redis 预扣库存(
DECR原子操作) - 异步落库到 MySQL(批量 INSERT)
- 超时未支付的库存回补(延迟消息或定时任务)
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 的,金融/积分场景必选。