资讯详情

资讯详情

实时事件聚合引擎REA:从架构设计到内存控制的完整实践

REA 是我最近完成的一个内部项目的代号全称是 Real-time Event Aggregation翻译过来就是实时事件聚合。做这个事之前我给业务团队写报表和指标查询支持每天对着几十张表、上百条过滤条件最痛的不是采集而是不同系统各自算一套口径永远对不上。我当时的想法很简单写一个足够轻的实时聚合服务把事件收进来按规则算出结果再提供统一查询。折腾了大概两个月如今回头写这篇复盘。这个项目到底适合谁看如果你平时要处理行为日志、接口监控指标、业务事件计数尤其是那种“每秒钟都有大量小事件进来但要等到几秒后才能看到整体数字”的场景那这篇的经验可以直接抄。反过来如果你的数据量只有每秒几千条随便写个脚本轮询数据库就够用了REA 这套架构确实有点重。但架构思想是通用的哪怕你只是做一个小工具后面讲到的窗口设计、哈希聚合、内存控制思路也都能用。1. 为什么自研一个叫“REA”的聚合服务1.1 最初一个月我到底在调通哪些数据链路项目源头看起来简单不同服务每天产生事件日志格式五花八门有的是 JSON 行有的是普通文本有的走消息队列有的直接往日志文件里写。我需要统一个入口让这些数据进来之后能按分钟、按小时、按用户维度算出结果。最初我并没有想自己造轮子而是把精力花在“对齐口径”上结果发现这才是最耗时的。举几个具体场景接口耗时统计每条请求上报耗时需要按 API 分组算 P50、P95、P99。用户流量累计用户播放或下载行为要按天累加用量逾期清零。活跃去重统计某个页面的独立访客数要求去重精度尽量高。这三个场景看起来互不相同但本质都是“把一条一条的事件按某种规则聚合起来”。我最初用脚本在各处跑批逻辑散落在不同仓库里某天指标对不上就要翻三四个项目的代码。于是我想做一个统一服务把这些聚合能力集中管理。需求拆到最后只剩一句话实时接收事件按用户配置做聚合对外提供查询接口。它不需要很强的事务能力也不需要存全量明细只要把结果给出来哪怕有几秒延迟也能接受。但延迟不能太离谱比如实时看板就别等到 10 分钟以后才出数。1.2 看了几套方案之后我发现最缺的是“轻”做技术选型时我认真对比了几类方案。第一类是重量级流处理框架部署一套集群至少要 3 个节点配置文件和运维成本都不低第二类是商业 SaaS接口好用但数据出了公司边界很多团队不敢用第三类是现成的时序数据库写入和查询都不错但维度多、乱序多的时候灵活性和资源占用是个问题。我用一个表格把这几个选项列出来方案部署成本实时性灵活性运维负担自研 REA低单进程即可高秒级出数高规则可自定义低重量级流处理框架高需要专门集群高中高商业 SaaS低开箱即用中低受接口限制无但数据在外部时序数据库中依赖外部组件中中中这里不是要否定那几类方案它们在大规模场景下都很好但我的目标场景非常明确单机就能跑、数据不出内部网络、聚合规则需要按业务场景随意调整。这么一筛自研一个精简服务反而是最省事的。项目代号定为 REA也正是因为它的核心就是事件聚合。内部的时候大家叫它“啊”现在对外复盘统一说 REA。2. 架构取舍一进一出中间全靠内存2.1 整体数据流采集、缓冲、聚合、查询四层REA 的架构没有花哨的分布式设计就是一个常驻内存的单进程服务。整体数据流分成四层接入层、缓冲层、聚合层、查询层。接入层负责接收事件。我一开始只做了 HTTP 接口客户端直接把 JSON 行 POST 过来。后来发现有的服务端对开销敏感又加了一个 UDP 端口用轻量级协议收事件。UDP 的好处是快坏处是可能丢包所以只建议用在容忍少量缺失的场景。缓冲层解决“背压”问题。如果聚合层处理不过来不能把 HTTP 请求卡死所以我用一个有缓冲的 channel 把事件暂存。聚合层消费速度跟不上时最坏的情况是丢弃或返回限流响应不能影响业务服务本身。聚合层是核心后面单独讲。查询层则是对外暴露 HTTP 接口读取聚合结果返回给调用方包括最近一分钟、上一完整分钟、指定时间段等不同粒度。四个层各司其职每一层都可以独立启停和扩展。最初我只写了聚合和查询两层后来压测发现 HTTP 接入和聚合耦合会让 GC 变重才把接入层单独拆出来。2.2 窗口模型滚动窗口为什么比自然语言好理解这类系统里最难讲清的就是“多久算一次”。有的数据需求写的是“最近 5 分钟”这其实是滑动窗口有的写的是“每分钟整点汇总”这是滚动窗口。REA 默认采用滚动窗口把时间切成固定长度的桶比如 60 秒一个桶。事件进来后按照事件发生时间落到对应桶里桶达到结束时间后就不再接收新数据进入可查询状态。滑动窗口虽然更贴近直觉但实现和查询都复杂很多。连续两个请求一个查“最近 5 分钟”另一个查“最近 1 分钟”存储和计算都要反复扫描。滚动窗口有一个巨大的好处窗口是确定的同一批事件无论查多少次结果都一样不会因为时间推移而改变。这是排障时最重要的一点。如果今天早上 10 点查出来的 P95 和下午 2 点查出来的 P95 不一致只可能是因为数据补录或者窗口边界算错而不是窗口数据还在变化。2.3 快照怎么落才不影响热路径性能纯内存存储速度快但进程重启会丢数据不符合实际要求。我选择的是“热数据不落盘、冷窗口异步落盘”的策略。具体来说只有已经结束并进入可查询状态的窗口才会被写到磁盘正在活跃写入的当前窗口始终留在内存。后台有一个 goroutine 每隔 5 秒扫一遍窗口状态把超过 30 秒的窗口固化到本地文件。查询的时候当前窗口直接读内存历史窗口则先查内存缓存没有再走磁盘加载。这样既保证了实时性又不会因为每次聚合都写磁盘把性能拖垮。第一次压测时我把快照写成同步的边聚合边写文件。结果 QPS 一到 3 万磁盘 I/O 就把整机拖慢了。改成异步批量写入后性能马上就稳了。这个教训让我深刻意识到聚合系统里热路径和持久化路径一定要分开。3. 核心实现哈希聚合、滑窗与内存控制3.1 聚合桶的数据结构聚合层本质上是一个“维度键到指标值”的多层哈希表。以“按 API 分组统计耗时”为例维度键是 API 名称指标是 count、sum、P95。REA 用 map[string]MetricBucket 作为核心结构每个维度键对应一个小的聚合对象。指标类型我实现了三种Counter数值累加适合计费、次数统计。Gauge最新值覆盖适合实时温度、连接数。分位数估算用近似算法保存数据分布适合耗时、延迟这类指标。分位数是最麻烦的部分。若要精确算 P95需要保留所有样本或很粗的分桶。REA 默认用了一个固定大小的分位数采样器每个采样器内部保留最多 1024 个样本点插入时随机替换旧样本查询时按位置估算。这种近似误差在 1% 左右对监控场景完全够用。下面是核心结构的伪代码type MetricBucket struct { Count int64 Sum float64 Last float64 Quantile *StreamEstimator // 近似分位数 } type Aggregator struct { shards [64]Shard } type Shard struct { mu sync.Mutex buckets map[int64]map[string]*MetricBucket // 窗口时间戳 - 维度 - 指标 }从查询角度看所有数据最终都在内存里所以响应时间基本是微秒级。维度键越少聚合越快维度键膨胀到千万级别内存就会告急因此 3.3 节讲的控制手段必不可少。3.2 分钟级滚动窗口轮转的实现滚动窗口实现我采用的是“两级切换”。外层是一个环形数组数组长度按保留时长决定比如保留 2 小时就分配 120 个桶位。每个桶位对应一个时间戳表示这一分钟窗口的开始时间。有一把后台定时器每秒检查一次当前时间如果距离当前桶的结束时间还有 1 秒就预创建下一个桶并把上一秒结束的窗口标记为“可查询”。这个预创建非常重要否则整点切换瞬间所有写请求都会挤到新桶创建逻辑里造成延迟抖动。事件进入时先用事件时间戳除以窗口长度得到窗口索引。如果索引超出了保留范围就丢弃如果落在当前窗口或上一窗口就正常合入更旧的窗口因为已经固化到磁盘不再接收事件。这样实现有两个好处窗口切换不需要遍历全量数据。查询引擎可以只读两个窗口当前窗口和上一完整窗口。3.3 内存不失控的三个约束内存是 REA 最容易被攻击的点。事件量一大维度一多分位数采样器再吃一点内存涨起来很快。我总结了三个必做的约束。第一是维度键数量上限。每个窗口里允许的维度键总数有一个硬上限默认 20 万个。超过之后不再新建维度键而是把这部分事件计入一个无意义的“overflow”维度。这样可以防止恶意或异常事件把内存打爆。第二是事件时间偏移窗口限制。如果事件里的时间戳比当前时间晚很多或者早很多都不能进入实时聚合。REA 只接受“最近 5 分钟以内”和“未来 30 秒以内”的事件。这样能挡住乱序太离谱的数据。第三是采样器大小自控。分位数采样器不是无限增长而是固定 1024 个点满了之后用“随机替换”策略。这意味着每个维度键的额外内存是可控的不会因为某个极端热点一直涨。我第一次上线时没做维度键上限结果一个测试脚本用随机维度名灌了半小时进程内存从 400MB 涨到 6GB最后被系统杀掉。补齐上面三个约束后内存曲线再也不出现那种无限增长的状态。4. 从配置到代码REA 怎么落地4.1 启动配置与核心参数REA 的默认配置尽量简化但有两个参数必须认真填窗口大小和维度键上限。窗口大小决定出数粒度维度键上限决定内存预算。其余比如端口、保留时长、快照路径都给了合理默认值。下面是一份 YAML 配置示例server: http_port: 8081 udp_port: 9001 query_port: 8082 window: size_seconds: 60 retention_minutes: 120 event_time_max_skew: 5m aggregation: max_dimensions_per_window: 200000 quantile_samples: 1024 flush_interval_seconds: 5 storage: snapshot_dir: ./snap compress: true实际使用时我把查询端口和写入端口分开是为了避免读写争抢同一个 HTTP listener。查询流量通常是低频的但偶发的大查询会把写入连接拖住端口隔离之后写入路径不会因为查询慢而积压。事件格式是我自定义的一个高性能 JSON 子集核心字段包括ts、dimensions、metrics、type。举个例子{ts: 1712345670, dim: api:login, m: {cnt: 1, cost: 23.5}}这样的结构简单序列化和反序列化成本很低。如果业务字段特别多可以配置字段白名单不参与聚合的字段直接忽略避免解析开销。4.2 核心聚合逻辑参考实现聚合入口的代码尽量精简流程是解析时间戳检查窗口归属找到对应分片的锁更新指标。我写的 Go 参考实现大约不到 300 行下面贴出最关键的一段func (a *Aggregator) Add(event Event) error { if event.Timestamp a.currentWindowStart()-skewSeconds { return ErrTooOld } if event.Timestamp a.currentWindowStart()maxFutureSkew { return ErrTooFuture } windowTime : event.Timestamp - (event.Timestamp % windowSeconds) shard : a.shards[hashKey(event.Dimension)%shardCount] shard.mu.Lock() defer shard.mu.Unlock() bucketMap, ok : shard.buckets[windowTime] if !ok { if len(shard.buckets) maxBucketsPerShard { // 简单淘汰丢掉最久远的窗口 for k : range shard.buckets { if k windowTime-retentionSeconds { delete(shard.buckets, k) break } } } bucketMap make(map[string]*MetricBucket, 16) shard.buckets[windowTime] bucketMap } b, ok : bucketMap[event.Dimension] if !ok { if len(bucketMap) maxDimPerWindow { return ErrDimensionOverflow } b newMetricBucket() bucketMap[event.Dimension] b } b.Count event.Metric.Count b.Sum event.Metric.Value b.Last event.Metric.Value if b.Quantile ! nil { b.Quantile.Add(event.Metric.Value) } return nil }这段代码的核心在于加锁粒度。我用 64 个分片每个分片管理一批窗口和维度键。正常情况下不同维度的写入会分散到不同分片上并发冲突很少。查询时同样按维度 key 定位分片相当于把全局锁拆成了局部锁。如果你自己实现一个简化版不需要分片一把全局锁也能跑通小流量。但一旦压到 10 万级 QPS锁竞争会成为瓶颈。分片不是优化是必要设计。4.3 压测结果与参数调整我用的压测机器是 8 核 16GB 的普通开发机压测工具往 HTTP 接口发模拟事件事件内容是“随机 API 名加耗时”共 100 个 API压测场景模拟的是真实监控上报。场景写入 QPS查询 QPSP99 写入延迟P99 查询延迟内存占用单写入单查询6.2 万300038ms5ms650MB写入和查询同时压5.8 万600042ms8ms780MB大维度场景5 万维4.1 万200055ms12ms1.2GB从结果看写入 QPS 达到 6 万时瓶颈已经不在聚合逻辑而在 HTTP 请求解析。后来我把 HTTP 请求体改成读一行处理一行避免把整个 body 先读进内存写入延迟降低了差不多 10%。查询侧的优化更明显。最初查询都是直接扫描 map维度多时 CPU 飙高。后来给每个窗口的维度键加了一个lastAccessed时间戳查询时优先读最近访问过的维度热门 key 基本都在缓存路径上冷 key 查询才会落慢路径。这个改动很小但收益非常直观。5. 上线后踩过的坑重复、乱序和热点5.1 重复数据一次性事件被客户端重试了上线后第一个星期某个业务方反馈“播放量对不上”。排查后发现他们的服务端在调用 REA 接口时如果超时就会重试而 REA 已经把第一次请求聚合进去了重试又加了一次计数直接翻倍。这个问题在聚合系统里很难通过简单去重解决。REA 不是明细存储不可能把所有事件主键都记住。我采取的方案是让客户端使用幂等 IDREA 接收层用一个短期的布隆过滤器缓存最近 5 分钟的幂等 ID如果重复调度就直接丢弃。布隆过滤器有误判但只会把本来应该计入的少量事件误判为重复整体误差在万分之几业务可以接受。更重要的是这个过滤器是可重启恢复的重启后从特定日志重新加载最近 5 分钟的 ID解决了服务重启导致误判的问题。5.2 乱序事件窗口归属总是拿不准另一个坑来自事件时间戳。业务端上报用的可能是本地时间某个服务时钟偏了 3 分钟那它产生的事件在过去 3 分钟的窗口里已经固化REA 默认会丢弃结果那个服务的数据就一直偏小。我之前设定的时间容忍度是“最近 5 分钟以内”但固化窗口是从 30 秒开始的。也就是说一个事件如果延迟 40 秒到达它对应的窗口已经进入异步落盘但还没完全删除这时该不该合入我的最终处理窗口固化后 30 秒内允许补写并把窗口状态标记为“已更新待重新落盘”。超过 30 秒就拒绝。这样避免了抖动窗口导致查询结果反复变化也允许了少量延迟数据。5.3 热点 key 拖垮单节点第三个问题是最经典的“热点 key”。某个营销活动上线后所有事件都汇聚到同一个 API 维度上按 hash 定位到的那个分片 CPU 被打满其他分片却很闲。修复方案不是改锁而是“拆维”。在内部给每个维度 key 后面追加一个分片号比如api:login变成api:login#0到api:login#15。查询时再把这些子维度合并。这样热点维度被拆成 16 份均匀分布到 16 个分片上单分片压力直接降下来。代价是查询变复杂聚合时需要多扫 16 个子维度再合并。但热点场景必须做均衡否则单点故障随时可能出现。5.4 一次内存突刺的排查有一天监控显示内存从 800MB 跳到 2GB持续 20 秒后又回落。我看日志发现是整点窗口切换时大量窗口在同一秒被创建然后又立刻填了大量事件内存分配器来不及回收之前的旧窗口数据。优化方式有两个。第一个是把窗口预创建提前到当前窗口前 20 秒而不是最后 1 秒摊平创建开销。第二个是内存池复用旧窗口指标对象不直接释放而是清空后交给对象池下一个窗口复用。这样整点切换时的分配压力大幅减少。内存突刺是这类系统最常见的隐蔽问题从监控曲线看就是一秒之间的“尖刺”。把分配行为变成可预测的往往比追求最优算法更有效。6. 如果重做一次我会改掉的三件事REA 走到现在核心功能已经完全满足需求。但如果让我回到最初有三个方面一定会重新设计也算是我最真实的经验。第一事件协议会提前定一个 Schema。最初为了快速接入每种事件类型都自带字段聚合逻辑写了很多分支。后来才发现统一协议虽然是约束但反而让接入方更容易理解。字段白名单加上通用dimensions和metrics比“先自由写再靠文档对齐”要靠谱得多。第二不会一开始就想“分布式”。我花了一些时间考虑横向扩展、分片调度、数据同步最后发现单机模式跑了很久都没到瓶颈。过早预估规模是浪费。正确做法是先做单机稳定版本等单机上限不够时再考虑多机而不是开局就背上复杂设计。第三会把查询能力和聚合能力分开成两个独立服务。现在 REA 聚合和查询在同一个进程短期方便但查询侧如果引入复杂表达式必然会拖慢聚合侧。如果拆开查询服务可以单独扩机器这种解耦在数据量大以后价值会越来越明显。把这三个点记录下来不是劝大家也用 REA 这个名字而是想说很多看似“没有技术含量”的后端服务真正做完以后才知道坑都在细节里。实时系统没有玄学大多数问题就是窗口、锁、内存和重试这四个维度的权衡。希望这篇复盘能让你在写自己的聚合服务时少走一段弯路。
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →