ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

基于滑动窗口去重设计轻量实时事件聚合器:从需求到落地

基于滑动窗口去重设计轻量实时事件聚合器:从需求到落地 之前我负责的一个内部工具项目代号就叫 rea全称后来定为 Realtime Event Aggregator。最初催生这个项目的需求并不复杂多条业务线不停地产生事件每秒钟少则几千、多则几万条我们要把这些事件统一收上来按业务维度做实时聚合算出的指标要能马上给数据看板和大屏用。看起来像是现成流处理框架能干的活但受限于部署环境和团队技术栈最后决定自己写一个轻量、可控的聚合器于是就有了 rea 这个项目。这篇文章不打算讲特别宏大的架构只把 rea 从需求拆解、模块设计到核心代码实现以及我在实际部署和维护中踩过的坑原原本本梳理出来。如果你正准备做类似的实时事件处理系统这套设计和踩坑记录可以直接当作参考。1. 项目概述为什么要自己造一个 rea1.1 最初的需求场景我接手项目时业务方给的需求是“把分布在多个服务和客户端的行为事件统一收集起来按频道、按类型、按时间维度聚合最终产出实时排行、趋势曲线和异常告警”。比如用户点击、支付成功、播放完成、设备心跳这些事件格式千差万别来源系统也不一样。有一部分会直接写入消息队列有一部分会打到内部网关还有一部分是客户端批量上报。如果不做任何加工直接把这些事件往存储里灌问题很大。首先是数据量高峰期每秒几万条写入压力会直接打爆存储其次是口径不统一有的系统用时间戳毫秒有的用秒有的根本没传时间字段下游算同比环比时根本对不上再一个是重复事件客户端断网重试、消息队列重投都会导致同一条业务事件在链路中被重复消费。所以整个项目的核心目标非常明确在事件被下游消费之前先把脏数据清理掉把不同来源的事件变成统一格式再按业务需要的维度做滑动窗口聚合产出的结果必须可解释、可回溯。1.2 rea 的定位与目标rea 定位不是一个大而全的流处理平台而是一个“介于接入和计算之间的轻量聚合服务”。它要满足四个硬性要求多协议接入至少要支持 HTTP 上报、消息队列消费、以及公司内部网关的转发统一事件模型不管来源是什么进入 rea 后必须遵循同一个事件结构实时聚合计算以秒级窗口为粒度输出计数、去重数、均值、分位数等指标故障可恢复服务重启后未处理的缓冲数据不能全部丢光至少要能恢复到最近一个存档点。在早期设计时我参考过一些开源流处理框架也纠结过要不要直接用现成的。但当时团队对框架的运维经验有限而且业务模型比较特殊事件不是简单的 key-value而是带层级结构的嵌套对象很多聚合逻辑还依赖前一次聚合的结果做增量计算。用通用框架反而要把业务逻辑硬绕成 SQL 和 UDF开发和排障成本都不低。最终决定自己维护一个核心不超过三千行的聚合引擎rea 就这样定下来了。2. 整体设计拆解从需求到架构2.1 统一事件模型rea 内部的每一次上报都先被包装成一个 ReaEvent。这个结构看起来简单但设计它的时候踩了不少坑。核心字段如下eventId全局唯一的事件 ID用于去重source来源标识比如“app_ios”“order_server”eventType事件类型比如“click”“pay_success”“heartbeat”channel业务频道比如“首页推荐”“直播弹幕”occurredAt实际发生时间毫秒级时间戳receivedAtrea 收到事件的时间毫秒级时间戳payload业务附加数据JSON 对象长度不限。一开始我图省事直接让上游把原始 JSON 原样传进来我在内部解析。后来发现不同来源的时间字段名完全不一样有的是time有的是ts有的是create_time每次兼容一个上游都觉得在给补丁打补丁。重置点所有事件进入 rea 的第一件事就是做字段映射统一成上面的结构。凡是映射失败的事件不会直接丢弃而是打入“脏数据队列”方便回头排查。这里有一个很重要但容易忽视的点occurredAt和receivedAt必须分开存。很多聚合场景看的是业务发生时间而链路延迟分析看的是接收时间两个字段对不上才能定位是上报延迟还是处理性能问题。2.2 管道处理流程rea 的处理主链路是一条管道核心环节只有四个接入、标准化、聚合、输出。每个环节之间用有界队列缓冲。接入层Ingestion接收来自 HTTP 接口或消息队列的事件这个阶段只做最基础的合法性校验防止空报文和畸形 JSON 打穿后续逻辑标准化层Normalizer把来源不同、字段不同的原始事件映射为统一的 ReaEvent 格式聚合层Aggregator负责时间窗口内的计数、去重、均值计算窗口结束后输出聚合结果输出层Exporter把聚合结果写入下游比如 Redis、ClickHouse、消息队列或者直接通过 WebSocket 推给大屏。需要注意这个管道模型不是简单的一个线程处理完再交给下一个线程那样会让整体吞吐量受限于最慢环节。实际实现中每个环节都有独立的线程池环节之间通过队列解耦。接入层如果瞬时流量过大不会把压力直接传给聚合层而是先在队列中堆积给聚合层留出削峰的时间。2.3 为什么不用现成流处理框架这个决策值得单独说说。当时我也做过对比现成框架的最大优势是生态完善窗口、状态、容错都现成。但如果只是做“分钟内实时统计”这种量级引入一套完整框架会带来三个额外成本部署成本框架通常需要集群协调组件运维复杂度上升一个台阶学习成本团队成员要重新学习编程模型排障链路变长业务耦合很多聚合逻辑是对业务指标的特化用通用框架表达并不优雅。rea 的通用性虽然不如大框架但对内部场景来说足够用。后续如果需要扩容也只需要把聚合层拆成多实例加上一层分区协商就行。这也是我坚持“轻量优先”的主要原因。3. 核心模块实现细节3.1 接入层多源事件消息的接收接入层的第一版实现只做了 HTTP JSON 上报。后来发现消息队列里的存量事件也要消费就把接入层做成了可插拔的接口核心代码如下public interface EventSource { void start(EventHandler handler); void stop(); }HTTP 接入用一个轻量 HTTP 服务实现每收到一条 JSON 就转成内部的RawEvent然后写入接入队列。消息队列接入则使用消费组模式循环拉取消息并调用同一个 handler。这里要特别注意消息队列的消费确认时机。最开始我在处理完事件后才提交 offset导致聚合引擎反压严重时消费线程会一直阻塞最终触发消息积压。后来改成先提交 offset再异步处理但这样做又会增加事件丢失的风险。最终采取的方式是本地先做持久化缓冲确认数据写入本地缓冲文件后再提交 offset这样既保证了不积压又能在服务崩溃后重新加载未处理数据。3.2 聚合引擎时间窗口与去重聚合引擎是 rea 最核心的部分。我用的是滑动窗口加增量计算模型窗口默认 10 秒滑动步长 2 秒。每个窗口维护一个指标状态对象记录当前窗口的事件量、去重数量、以及各个维度的子计数。窗口的具体实现方式是“按开始时间分桶”。每来一条事件计算它所属窗口的开始时间然后找到对应的桶更新状态。窗口结束时把状态对象交给输出层同时创建新的窗口桶。这种实现比写一个 Timer 定时触发要简单得多天然支持乱序事件的归窗。去重逻辑是另一块容易踩坑的点。最早的实现是直接在内存里维护一个ConcurrentHashMapkey 是事件 IDvalue 是时间戳。这种做法在小流量下没问题但每秒几万条事件时内存会被大量只有几 KB 的键值对占满GC 压力也非常大。后来我改成了 BloomFilter 加固定长度环形缓冲的组合。BloomFilter 负责回答“这个 eventId 以前是否见过”环形缓冲负责保存最近 N 秒的 eventId 全集用于窗口结束时生成精确去重数。BloomFilter 会出现误判也就是“没见过但实际上就是新的”这种情况几乎不会发生但存在小概率把新事件误判为重复事件。为了避免业务上不可接受我在 BloomFilter 前置了一个精确去重缓存只有当事件时间戳落在当前活跃窗口范围内时才走精确判断。这样既保证了内存可控也保证了窗口内去重结果完全准确。3.3 存储与查询设计聚合结果不能只存在内存里必须落盘到下游便于查询。rea 默认输出两条链路近实时链路每秒把窗口聚合结果写入 Rediskey 格式为rea:{channel}:{eventType}:{windowStart}value 是一个 JSON 结构。大屏和告警直接读取 Redis离线链路每五分钟把窗口结果批量写入存储。因为离线链路的写入频率低批量大所以采用累积批量写入模式避免每窗口一次连接。需要注意的是输出链路要做失败保护。如果 Redis 短暂不可用不能把聚合线程阻塞太久。我在 Exporter 中加入了带超时的异步队列写入失败的重试次数最多三次三次以后直接丢弃并打日志。这种方式丢了数据但至少保住了主链路。4. 实操过程从零搭建 rea 核心链路4.1 依赖选型与版本因为 rea 是纯 Java 项目核心依赖并不多组件用途版本说明NettyHTTP 接入4.1.x自研消息队列客户端消费上游事件内部版本Caffeine本地缓存/去重3.xGuavaBloomFilter31.xJacksonJSON 解析2.15.xRedissonRedis 客户端3.20.x选 Caffeine 而不是自己写 ConcurrentHashMap 缓存是因为 Caffeine 支持基于时间的过期策略正好可以配合窗口数据清理。Guava 的 BloomFilter 在内存占用和误判率之间可以调参默认设置误判率 1%事件量在每秒两万条以下时内存占用可以控制在 500MB 以内。4.2 关键代码实现聚合引擎的核心就是一个窗口管理器我用下面的代码做骨架说明public class SlidingWindowAggregator { private final long windowSizeMs; private final long slideSizeMs; private final ConcurrentHashMapLong, WindowBucket buckets new ConcurrentHashMap(); private final EventDeduplicator deduplicator; public SlidingWindowAggregator(long windowSizeMs, long slideSizeMs, EventDeduplicator deduplicator) { this.windowSizeMs windowSizeMs; this.slideSizeMs slideSizeMs; this.deduplicator deduplicator; } public void addEvent(ReaEvent event) { long windowStart event.getOccurredAt() - (event.getOccurredAt() % windowSizeMs); // 这里要注意滑动窗口需要维护多个重叠窗口分片 long currentWindowStart alignToSlide(windowStart); WindowBucket bucket buckets.computeIfAbsent(currentWindowStart, k - new WindowBucket(windowSizeMs, slideSizeMs)); if (deduplicator.isDuplicate(event.getEventId())) { bucket.incrementDuplicate(); return; } bucket.add(event); } private long alignToSlide(long start) { return start - (start % slideSizeMs); } }这段代码不是完整实现但展示了几个关键点。第一WindowBucket内部维护一个环形数组数组长度等于windowSizeMs / slideSizeMs每个元素是一个子窗口。事件到达时根据时间分布写入对应的子窗口。第二窗口结束不是靠定时任务而是在某事件的时间戳越过当前窗口边界时触发这样能自动处理空闲窗口避免大量空窗口占用内存。最终聚合结果的触发代码如下public ListAggregateResult flushExpiredWindows(long now) { ListAggregateResult results new ArrayList(); for (Map.EntryLong, WindowBucket entry : buckets.entrySet()) { WindowBucket bucket entry.getValue(); if (bucket.isExpired(now)) { results.add(bucket.toAggregateResult()); buckets.remove(entry.getKey()); } } return results; }flush 操作由后台线程每隔 500ms 调用一次所以下游拿到结果的延迟最多是窗口大小加 flush 间隔时间。4.3 参数计算与配置参数值不能靠拍脑袋定下面几步是我实测下来比较靠谱的方式。窗口大小业务上希望看到“最近 10 秒的点击趋势”所以窗口定为 10 秒滑动步长大屏每 2 秒刷新一次所以滑动步长定为 2 秒确保每次刷新能拿到新结果内存预算按每秒 2 万事件、每个事件 200 字节估算缓冲区内存约 4MB/s窗口内需要保存 5 个滑动分片每个分片又保存最近一个窗口的事件明细所以计算上瞬时内存峰值要按事件量 * 窗口大小 * 平均事件大小来估算。比如每秒 2 万、窗口 10 秒内存峰值就是 2 万 * 10 * 256 字节 ≈ 48MB再加上去重 BloomFilter 和缓存控制在 500MB 内完全可行。还有一个参数容易忽视BloomFilter 的大小。我按每日去重事件量 500 万来计算误判率 1%需要约 4800 万个 bit占用内存约 60MB这完全可接受。但如果误判率调成 0.1%内存会翻到 90MB 左右需要根据线上内存水位决定。配置参数最终会写在 YAML 文件里rea: window: size-ms: 10000 slide-ms: 2000 flush-interval-ms: 500 dedup: bloomfilter: expected-insertions: 5000000 fpp: 0.01 cache-expire-ms: 15000 ingest: http-port: 8080 batch-size: 200 exporter: redis-key-prefix: rea max-retry: 35. 常见问题与排查实录5.1 窗口数据重复或丢失实际运行中遇到最多的问题是“同一窗口的数据被重复输出”。最初我在输出层没有做幂等聚合结果写入 Redis 是直接SET覆盖看似不会重复。但当输出线程重试时重试用的还是同一份 AggregateResult 对象第二次写入就会覆盖第一次的结果看起来倒没产生重复值。真正产生重复的情况来自下游消费如果下游从消息队列消费聚合结果而 rea 的消费语义是 at-least-once下游没做去重就会看到两条相同窗口结果。解决方式有两种。一种是在输出消息里带上windowStart和aggregateKey下游用这两个字段做唯一键重复消息直接忽略另一种是 rea 自己保证每条窗口结果的输出只有在“窗口确认结束”后才能发出去。我最终两种都做了上游幂等下游兜底。5.2 事件乱序导致聚合偏差事件发生时间和到达时间往往不一致客户端离线缓存会导致一批事件延迟到达。如果严格按occurredAt归窗那么延迟数据会重新计算历史窗口已经输出的结果就没办法修正。如果严格按receivedAt归窗那么离线事件会被归到错误的时间段趋势曲线会出现“时间漂移”。我的妥协策略是“双时间戳校验”。事件默认按occurredAt归窗但如果receivedAt - occurredAt超过窗口大小的两倍说明事件延迟严重直接归入一个“迟到了但还是要算”的补偿队列在下一分钟输出时对上一分钟的目标指标做一次性修正。这样做不追求绝对准确但不会让用户在看大屏时突然看到旧数据把曲线拉歪。5.3 高水位与内存泄漏排查rea 部署一段时间后发现内存缓慢上涨最终 OOM。一开始以为是窗口桶没有释放后来排查发现是 BloomFilter 的缓存填满后没有按时间清理。我用 Caffeine 缓存 eventId 时设置了 15 秒过期但事件量大的情况下Caffeine 的异步清理会有延迟内存峰值会比预期高。解决方式是在 flushExpiredWindows 的后台线程里主动调用一次缓存清理同时在 BloomFilter 的位数组上做定期重置。更好用的技巧是对 BloomFilter 做“双层滚动”每 5 分钟新建一个 BloomFilter旧的在 5 分钟后自动废弃这样既能保留最近事件的去重信息又避免整个位数组无限增长。5.4 排查速查表现象可能原因处理方式聚合结果延迟超过窗口数倍输出队列阻塞、Redis 线程池打满检查 exporter 线程池活跃数增加输出队列容量或接入备份链路事件丢失接入层反压丢弃开启 persists-buffer落地本地临时文件检查消费者 offset 提交是否过迟重复率异常高客户端重复上报、消息队列重投检查 eventId 生成唯一性确认 BloomFilter 误判率配置窗口结果偶尔跳变延迟事件补偿修正调整迟到的判定阈值确认补偿队列是否堵塞频繁 Full GC窗口内存超预算调大 flush 间隔把窗口分片数调小降低事件明细保留比例6. 最后再分享一点实际体会rea 这个项目做下来我最大的感受是实时聚合系统的复杂度不在于“实时”而在于“准确”。窗口怎么开、去重怎么做、延迟事件怎么处理每个环节的取舍都会直接影响最终指标的含义。不要盲目追求框架的大而全先把核心路径跑通再在迭代中逐步补全容错和运维能力这才是中小团队做实时系统最稳妥的路径。另外有一个小技巧想分享给做类似项目的同学聚合结果输出前一定给对方一个唯一的“结果ID”这个 ID 可以带上窗口时间和聚合维度 hash。下游不管是写库还是做缓存都拿这个 ID 做幂等。这个习惯帮你省掉后面非常多的对账工作也是我在 rea 上线后处理重复消费问题时觉得最值的一个设计。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进