用 Project Reactor 的 bufferTimeout 代替 BufferTrigger 做高并发计数聚合
本文最后更新于 2026-07-11,文章内容可能已经过时。
用 Project Reactor 的 bufferTimeout 代替 BufferTrigger 做高并发计数聚合
记录一次在计数服务里做「粉丝数聚合写」的实现。原教程用快手的 BufferTrigger,我改用了 Spring Boot 自带的 Project Reactor,零新增依赖。
背景:为什么要聚合
粉丝数是典型的高并发写场景——某个用户突然爆火,短时间涌入大量关注。如果每来一条 MQ 就 HINCRBY Redis + 写一次库,Redis 和数据库都顶不住。
聚合的本质是:把一个时间窗口内的大量 +1/-1 先在内存里合并,只对 Redis/DB 做一次净增量操作。
比如同一秒内,对用户 27 的粉丝数操作是:
+1 +1 -1 +1
无论顺序如何,聚合到一起就是净 +2,只需操作一次 Redis 与数据库。这就是关注/取关这类计数天然适合聚合的原因——它满足交换律,聚合无需关心顺序,最终一致即可。
聚合需要两个触发条件(谁先到谁触发):
- 数量阈值:攒够 1000 条,聚合一次;
- 时间窗口:即使 1 秒内不足 1000 条,为了数据及时性,到点也触发一次。
原方案:快手 BufferTrigger
原教程用的是快手开源的 com.github.phantomthief:buffer-trigger:
private BufferTrigger<String> bufferTrigger = BufferTrigger.<String>batchBlocking()
.bufferSize(50000) // 缓存队列最大容量
.batchSize(1000) // 一批最多聚合 1000 条
.linger(Duration.ofSeconds(1)) // 多久聚合一次
.setConsumerEx(this::consumeMessage)
.build();
// onMessage 里:
bufferTrigger.enqueue(body);
它内部就是一个带「数量阈值 + 时间窗口」双触发的缓冲队列。功能没问题,但这个库 2021 年(0.2.21)后就停更了,是单维护者项目。给一个全新模块引入一个多年不更新的第三方依赖不太理想,所以我换掉了它。
替代方案:Project Reactor 的 Sinks + bufferTimeout
Reactor 随 Spring Boot 传递依赖,已经在 classpath 上,零新增依赖,且由 Spring 团队活跃维护,原生支持背压。核心是三件套。
① 一个 Sink 当入口缓冲
private final Sinks.Many<String> sink = Sinks.many().unicast().onBackpressureBuffer();
Sinks.Many 是 Reactor 里「手动往响应式流里塞元素」的入口。unicast() = 单订阅者,onBackpressureBuffer() = 下游处理不过来时把元素缓存在内部队列里(而不是丢弃或报错)。
② 启动时订阅,挂上 bufferTimeout
@PostConstruct
public void init() {
subscription = sink.asFlux()
.bufferTimeout(1000, Duration.ofSeconds(1)) // ← BufferTrigger 的双触发等价物
.subscribe(this::consumeBatch);
}
bufferTimeout(1000, 1s) 就是等价物:攒够 1000 个元素,或距上批满 1 秒,谁先到就把这一批 List 推给下游。
发 3200 条就会切成 1000 + 1000 + 1000 + 200 四批——前三批靠「满 1000」触发,最后 200 条靠「1 秒超时」触发。
③ RocketMQ 消费回调只负责入队
@Override
public void onMessage(String body) {
sink.emitNext(body, Sinks.EmitFailureHandler.busyLooping(Duration.ofSeconds(1)));
}
onMessage 不做任何业务,只把消息塞进 sink 后立刻返回。聚合和落库在 Reactor 的另一条线程上异步进行。
关键坑:并发 emit 不安全
这是整个替代方案里最容易踩错的地方。
RocketMQ 是多线程回调 onMessage 的(日志里 ConsumeMessageThread_..._3 那个 _3 就是线程编号),但 Reactor 的 Sinks emit 并不是并发安全的。如果用 tryEmitNext,多个线程同时 emit 会返回 FAIL_NON_SERIALIZED 失败,消息就丢了。
所以这里不能用 tryEmitNext,要用带失败处理器的 emitNext:
sink.emitNext(body, Sinks.EmitFailureHandler.busyLooping(Duration.ofSeconds(1)));
busyLooping 的语义是:遇到并发竞争失败(FAIL_NON_SERIALIZED)时,自旋重试最多 1 秒,直到成功塞进去。这样多线程回调下也不会丢消息。
这一条是相比 BufferTrigger 唯一需要额外操心的地方——BufferTrigger 的
enqueue内部自己处理了并发。
聚合触发后做什么
拿到一批 List<String> 后的处理流程:
1. 逐条 JsonUtils.parseObject → List<CountFollowUnfollowMqDTO>(过滤 null)
2. FansCountAggregator.aggregate(dtoList)
→ 按 targetUserId 分组,FOLLOW +1 / UNFOLLOW -1 净算
→ Map<Long, Integer>,如 {68001: 1, 27: 3200}
3. 遍历 Map:仅当 redisTemplate.hasKey(key) 时才 HINCRBY fansTotal(不初始化缓存)
4. 把整个 countMap 的 JSON 再发一条 MQ 到 CountFans2DBTopic
→ 由落库消费者用 Guava 令牌桶(5000/s)削峰后 upsert 落库
聚合消费者主体代码:
private void doConsumeBatch(List<String> bodyList) {
log.info("## 聚合粉丝数消息, size: {}", bodyList.size());
// List<String> → List<CountFollowUnfollowMqDTO>
List<CountFollowUnfollowMqDTO> dtoList = bodyList.stream()
.map(body -> JsonUtils.parseObject(body, CountFollowUnfollowMqDTO.class))
.filter(Objects::nonNull)
.toList();
// 按目标用户分组净算增量
Map<Long, Integer> countMap = FansCountAggregator.aggregate(dtoList);
if (countMap.isEmpty()) {
return;
}
// 更新 Redis(仅当 Hash key 已存在,不初始化缓存)
countMap.forEach((targetUserId, delta) -> {
String redisKey = RedisKeyConstants.buildCountUserKey(targetUserId);
if (Boolean.TRUE.equals(redisTemplate.hasKey(redisKey))) {
redisTemplate.opsForHash().increment(redisKey, RedisKeyConstants.FIELD_FANS_TOTAL, delta);
}
});
// 转发落库 MQ,payload 为聚合后的 countMap JSON
Message<String> message = MessageBuilder.withPayload(JsonUtils.toJsonString(countMap)).build();
rocketMQTemplate.asyncSend(MQConstants.TOPIC_COUNT_FANS_2_DB, message, /* SendCallback */);
}
一个可测性改进:第 2 步「分组净算」的逻辑我特意抽成了独立的纯函数 FansCountAggregator.aggregate(),它不依赖 Reactor / MQ / Spring,输入 List<DTO>、输出 Map<Long,Integer>,可以直接单元测试。原教程是把这段逻辑塞在消费者内部的,不好测。
public static Map<Long, Integer> aggregate(List<CountFollowUnfollowMqDTO> dtoList) {
Map<Long, Integer> countMap = new HashMap<>();
if (dtoList == null || dtoList.isEmpty()) {
return countMap;
}
for (CountFollowUnfollowMqDTO dto : dtoList) {
FollowUnfollowTypeEnum typeEnum = FollowUnfollowTypeEnum.valueOf(dto.getType());
if (Objects.isNull(typeEnum)) {
continue; // 非法 type,跳过
}
int delta = switch (typeEnum) {
case FOLLOW -> 1;
case UNFOLLOW -> -1;
};
countMap.merge(dto.getTargetUserId(), delta, Integer::sum);
}
return countMap;
}
两个可靠性设计
① consumeBatch 整体 try/catch —— 别让管道断流
private void consumeBatch(List<String> bodyList) {
try {
doConsumeBatch(bodyList);
} catch (Exception e) {
log.error("## 聚合粉丝数批处理失败(已吞掉,避免终止订阅), size: {}", bodyList.size(), e);
}
}
consumeBatch 是 subscribe 的 onNext 回调。如果异常外溢,Reactor 会触发 onError 让整个订阅永久终止——之后所有粉丝数消息都没人处理,只能重启进程恢复。所以这里宁可丢掉出错的这一批(记 error 日志,由计数对账兜底),也绝不能让管道断流。
实测踩过一次:预热 Redis 时手误把值写成了
"0;"(带分号),HINCRBY要求字段是纯整数,于是抛ERR hash value is not an integer。正是这层 try/catch 兜住了,订阅没死,其他用户的计数照常处理。
② @PreDestroy 优雅关闭 —— 别丢缓冲区
@PreDestroy
public void shutdown() {
// 补发完成信号,触发 bufferTimeout 同步 flush 最后一批
sink.emitComplete(Sinks.EmitFailureHandler.busyLooping(Duration.ofSeconds(1)));
if (subscription != null) {
subscription.dispose();
}
}
停机时补发完成信号,让 bufferTimeout 把缓冲区里还没满 1000 条 / 未到 1s 的最后一批同步 flush 出去(unicast sink 的完成信号在调用线程上同步传播),再释放订阅。避免正常重启 / 发版时丢掉缓冲中的计数。
方案对比
| 维度 | BufferTrigger | Reactor bufferTimeout |
|---|---|---|
| 依赖 | 新引入,2021 年停更 | Spring Boot 自带,零新增 |
| 双触发(数量/时间) | batchSize + linger |
bufferTimeout(n, duration) |
| 并发 enqueue | 内部处理好了 | 需 emitNext + busyLooping 手动保证 |
| 背压 | bufferSize 有上限 |
onBackpressureBuffer(默认无界,需注意) |
| 逻辑可测性 | 逻辑在消费者内,不好测 | 净算抽成纯函数,可独立单测 |
一个共同的固有局限
两种方案都是「先 ack MQ,再异步聚合」——onMessage 把消息塞进缓冲后就返回,RocketMQ 随即提交 offset。如果进程在缓冲的这批消息落到 Redis/DB 之前硬崩,这批消息会永久丢失(无重投),也就是「至多一次」语义。
@PreDestroy 只能覆盖优雅停机(重启、发版),覆盖不了 kill -9 / OOM / 断电这类硬崩。要彻底兜底,需要一个计数对账任务——定期用关系表(t_following / t_fans)重算真实计数,校正计数表。这一点无论用 BufferTrigger 还是 Reactor 都一样,是聚合写这个模式本身的取舍。
一个待优化点
onBackpressureBuffer() 默认是无界队列。极端情况下(下游 Redis/DB 卡死、消息暴涌),可能堆积占内存。BufferTrigger 的 bufferSize(50000) 是有上限的。如果要更贴近有界行为,可以给 sink 配一个容量上限。生产环境值得关注这一点。
小结
- 用
Sinks.many().unicast().onBackpressureBuffer()+bufferTimeout(1000, 1s)就能等价替代 BufferTrigger 的「数量/时间双触发聚合」,零新增依赖。 - 务必用
emitNext + busyLooping,不能用tryEmitNext——RocketMQ 多线程回调下 Sinks emit 非并发安全。 - consumeBatch 要整体 try/catch,否则一次异常会终止整个订阅、管道断流。
- 用
@PreDestroy + emitComplete优雅关闭,flush 掉缓冲区最后一批。 - 聚合的净算逻辑抽成纯函数,方便单测。
- 记住聚合写是「至多一次」,硬崩会丢缓冲——生产上要有对账兜底。