本文最后更新于 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);
    }
}

consumeBatchsubscribe 的 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 掉缓冲区最后一批。
  • 聚合的净算逻辑抽成纯函数,方便单测。
  • 记住聚合写是「至多一次」,硬崩会丢缓冲——生产上要有对账兜底。