ARTICLE DETAIL

资讯详情

深耕商务建站与企业官网运营的一线实战洞察。

自研轻量级实时事件分析系统:从埋点到秒级看板的实践与踩坑

自研轻量级实时事件分析系统:从埋点到秒级看板的实践与踩坑 去年做活动复盘的时候我对着后台的SQL报表愣了很久明明活动还在进行中业务方问当前实时新增有多少我答不上来。不是没有数据而是数据散落在日志、数据库和一堆临时脚本里等我把它们攒齐、清洗、汇总出来十分钟已经过去了。就是在那一天我决定自己动手做一套轻量级的实时事件分析系统代号就叫rea全称是Real-time Event Analytics。这套系统从设计到上线前前后后花了不到三周却彻底改变了我们处理现在到底发生了什么这类问题的方式。它没有做成多么庞大的平台只是一个能支撑内部运营、产品和开发同学实时看数的系统。这篇内容适合那些不想一上来就引入重型流计算框架的中小团队也适合想搞清楚埋点、采集、聚合、查询全链路到底是怎么回事的前后端工程师。1. 逼我动手的三个需求痛点1.1 活动大屏背后的实时其实是延迟统计我所在的团队负责一个中型产品日活不算夸张但运营活动非常频繁。每当首页改版或者营销活动上线业务方最关心的问题永远是现在到底有多少用户进来了点击率怎么样听起来很简单实际查询却要绕一大圈。事件日志被统一收在文件里每天的定时任务负责把它们载入数据库再生成一份离线报表。这意味着当天白天的数据要等到凌晨才能看到白天只能靠写临时SQL去查每次执行时间还取决于数据量。活动期间流量是平时的几倍临时查询一个不小心就是几十万行记录去聚合数据库CPU直接被顶满其他业务跟着遭殃。这种伪实时带来的痛点不只是慢。更关键的是它严重压缩了运营决策的反应时间。活动上线前半小时运营想根据实时数据调整入口位置或广告语等我们算出数字用户早就流失了。要解决问题我得让数据链路从小时级至少走到秒级。1.2 找现成工具时遇到的三个不合身着手之前我花了一周时间调研现有的开源和商业方案。结论是没有哪一款能直接塞进我们团队而不产生新的问题。第一类是商业数据分析SaaS。它们体验确实好接入也快但费用不低而且核心数据要传到第三方平台。当时我们产品里有一些核心业务行为的数据业务方明确不希望在外部系统留存这一条就把大部分SaaS方案排除了。第二类是开源重型组件。比如基于流计算生态的方案功能强大但需要配套的消息队列、计算集群、监控体系。我们团队一共就四五个人还要兼顾日常业务开发根本没有余力去维护一套分布式系统。为了每天几万个事件上流计算框架明显是杀鸡用牛刀。第三类是传统BI报表工具。它们擅长把已有表结构做可视化但对事件流这件事没什么概念。我想要的是从埋点采集到实时聚合到查询API一整条链路都能自己掌控的东西而不是在数据已经入库之后再手工建模。三种路数摆在一起差距很明显方案类型优点不合适的地方商业SaaS接入快、图表全费用高、数据出域、定制受限开源流计算全家桶扩展性强、生态完善运维复杂、对团队要求高传统BI可视化强、报表丰富不关心事件链路、实时性弱自研轻量系统完全可控、成本低需要自己维护、迭代1.3 用一句话给rea划清边界当时很多人劝我要不先用定时SQL凑合等数据量大了再说。我坚持要做是因为等数据量大往往永远等不到。关键不是数据量而是能不能在被问到时快速给出答案。所以rea的边界在一开始就定死了用一句话说就是一个轻量级的、秒级延迟的、只关注关键事件的内部实时事件分析系统。具体来说rea不做以下几件事不采集全量用户行为只收录我们关心的核心事件不做用户画像和广告投放决策不做毫秒级实时竞价那样的低延迟场景不追求分布式和高可用。这几条边界非常重要因为它们决定了技术上可以走简单粗暴但可靠的路线。边界划清之后整个系统的技术选型就变得非常轻松因为我只需要回答一个问题在每秒几千事件、查询频率不高的规模下怎么用最少的人力把链路打通下面是我当时的完整设计。2. 技术架构怎么在够用和未来扩展之间找平衡2.1 埋点SDK自己写还是用现成的埋点是一切分析的地基。调研了一圈现成的开源埋点SDK功能确实不少自动采集页面浏览、点击热图、用户属性全都有。但我掂量了一下决定自己写一个不到10KB的轻量SDK原因有三。首先是包体积。我们的前端页面本身挂了图表库、请求库随便一个商业埋点SDK动辄几十KB对移动端用户体验影响不是小事。自己写的SDK可以只保留事件发送、重试、节流三个能力。其次是数据协议的可控性。现成SDK的事件字段格式往往是固定的想加一个团队自定义的环境标识要么做二次开发要么在数据落库之后再清洗。自研SDK可以直接采用我们内部定义的事件协议采集端、传输端、存储端使用同一套结构省掉很多中间转换。最后是隐私合规。自研SDK可以明确地控制哪些数据允许采集例如默认不采集系统字体、屏幕亮度这类无意义但容易引起误会的字段。由于rea从一开始就面向内部场景我们不希望采集用户敏感信息自研可以让这个承诺写进代码而不是写进文档。这里给一个最简单的浏览器端发送逻辑function sendReaEvent(eventName, props {}) { if (!window._rea || !window._rea.userId) { return; // 未初始化或匿名场景 } const event { v: 1, // 事件协议版本 id: generateUUID(), // 全局唯一事件ID name: eventName, userId: window._rea.userId, anonymousId: window._rea.anonymousId, time: Date.now(), // 设备本地时间戳 props: JSON.stringify(props), }; if (navigator.sendBeacon window.location.protocol https:) { navigator.sendBeacon(/rea/track, new Blob([JSON.stringify(event)], { type: application/json })); } else { fetch(/rea/track, { method: POST, body: JSON.stringify(event), headers: { Content-Type: application/json }, keepalive: true, }).catch(() {}); } }为什么用sendBeacon因为它适合页面卸载场景下上报事件不阻塞页面跳转浏览器会在网络空闲时把数据发出去。如果没有这条很多点击后立刻跳走的事件就会白白丢失。2.2 消息管道别一上来就上Kafka消息队列是整个链路的缓冲层很多人第一反应是Kafka。Kafka确实好高吞吐、持久化、消费者组一应俱全。但对我们这个每天几十万事件、峰值每秒两三千请求的场景它的问题也很明显依赖独立集群运维心智高磁盘占用、副本配置、分区调优哪一项都要花时间。我的选择是先用NATS。这是一个极简的云原生消息系统安装一个二进制就能跑支持主题订阅吞吐量对rea的场景完全够用还自带持久化选项。团队没有专职运维NATS单节点崩溃了也能快速重启成本非常低。如果你不想引入新组件用Redis Stream也可以它一样能做到消息的持久化和消费组。我后来把两种都跑过压测结论在下面这张表里对比项NATSRedis StreamKafka安装复杂度低低依赖Redis高单机吞吐数万级/s数万级/s百万级/s持久化支持JetStream支持强运维成本低低高适合rea适合适合大材小用选择NATS之后我们用它的主题把事件从采集服务分发给聚合服务和归档服务。这样即使聚合服务重启事件也会留在JetStream中不会丢失。哪天流量真的涨上来了从NATS迁移到Kafka只需要改消费端的几行代码因为rea的业务逻辑全部在消费者内部与消息系统解耦得比较干净。2.3 存储选型PostgreSQL起步ClickHouse留后路事件明细最终要落到数据库里。最开始我的候选名单里有四个选手PostgreSQL、MySQL、ClickHouse、DuckDB。先排除DuckDB它是嵌入式分析型数据库做离线分析很爽但并发查询能力和多用户访问的支持弱一些不适合作为线上服务的主存储。MySQL是顺手就能用但数据分析场景下它的聚合能力确实不如PostgreSQL更何况后面我打算用预处理索引优化。ClickHouse是我心里的未来方向列式存储、压缩比高、聚合极快特别适合事件分析。之所以没有一开始就用是因为我们团队对它的运维经验不足同时ClickHouse更适合在数据量已经很大的情况下体现优势。几十万行数据在PostgreSQL里用一条带索引的SQL也能在几百毫秒内返回完全够用。所以在第一版里我选了PostgreSQL但留下了一个伏笔定义事件明细表时把event_time、userId、event_name这些分析字段单独作为事件维度宽表来建方便以后原样迁移到ClickHouse。表格长这样CREATE TABLE rea_events ( id VARCHAR(40) PRIMARY KEY, -- 事件ID幂等键 name VARCHAR(64) NOT NULL, -- 事件名 user_id VARCHAR(64), -- 用户ID哈希后 anonymous_id VARCHAR(64), occurred_at TIMESTAMPTZ NOT NULL, -- 事件发生时间 received_at TIMESTAMPTZ NOT NULL, -- 服务端接收时间 props JSONB NOT NULL DEFAULT {}, -- 扩展属性 created_at TIMESTAMPTZ DEFAULT NOW() ); CREATE INDEX idx_rea_events_name_time ON rea_events (name, occurred_at DESC);2.4 查询与展示用最少代码做内部看板存储定了之后展示层我犹豫了一下要不要接一个开源BI工具它们功能确实丰富可以做下钻、联动、权限管理。但我们的需求其实很单一几个关键指标的趋势、实时粗略计数、按事件名分组汇总。为一个单一需求引入一套完整的BI同样不划算。最后我用一个轻量API服务加一个不到200行的前端页面解决了。API服务负责查库、聚合、缓存前端页面就放几块图表实时趋势折线、事件排行榜、基础漏斗。数据格式统一用JSON返回前端用一套现成的图表库渲染。这个方案的好处是后续想换BI、想开放数据给其他系统只要API不变底层随便换。3. 核心实现rea从0到1的五个关键环节3.1 先定事件协议再写代码这是我认为整个项目最重要的决定。如果没有统一的事件协议后面每一个环节都会因为字段不匹配而加班。rea的事件协议在JSON层面就定义死了字段意义如下v协议版本整数从1开始。后续加字段就升版本消费端按版本做兼容。id事件唯一ID由前端生成UUID。这个字段为幂等去重而生后面踩坑部分会专门讲到。name事件名命名规则是对象_动作比如button_click、page_view、order_submit。userId和anonymousId登录用户ID经过哈希后的值和匿名用户ID用来做漏斗和留存。time设备本地时间毫秒时间戳。props扩展属性JSON对象允许不同事件带不同的业务属性。这里有个细节为什么不直接用Protobuf或者Avro因为rea是内部系统解析链路上的消费者只有采集服务和聚合服务JSON的解析开销完全不是瓶颈。Protobuf虽然节省带宽、类型约束强但多一层编译、多一层schema管理对三周内要上线的项目来说收益不抵成本。技术选型要放在具体约束下看脱离场景谈性能没有意义。3.2 采集服务批量写入和背压处理采集服务是前端埋点请求的第一站它的职责很简单接收事件、校验字段、写入消息管道。但简单不代表可以随意写。我见过不少同类系统采集接口一个一个往数据库插结果流量稍大就卡死。rea的做法是在服务内做两级缓冲。第一级缓冲是一个内存队列接收到的每一条事件先丢进队列由后台批量任务每500毫秒或者攒够1000条后一次性打包发给消息管道。第二级缓冲就是NATS本身它保证即使采集服务进程崩溃已提交到JetStream的事件也不会丢。缓冲带来的直接问题是怎么处理队列满了。如果生产速度远超消费能力继续往队列里塞只会导致内存溢出。我当时定了一个降级策略核心事件比如订单、支付结果不允许丢弃即使延迟也要保非核心事件比如普通的按钮点击在队列超过80%水位时直接采样丢弃一部分并打一条日志。写代码时大概是这个样子# 伪代码示意reactive处理的思路 def handle_event(event, is_criticalFalse): if not queue.offer(event, timeout_ms100): if is_critical: # 核心事件阻塞等待 queue.put(event) else: dropped_count.inc() logger.warning(queue full, drop non-critical event) def batch_flush(): while True: batch queue.take_n(max_count1000, timeout_ms500) if batch: nats.publish(rea.event, batch)背压处理是最容易被忽略的。很多自建采集系统都是从单条写入改成批量写入就完事了完全不考虑生产者太快会怎样。结果就是突发流量一来服务直接内存溢出。3.3 实时聚合计数器放Redis明细留给数据库实时看数和精确统计是两种不同的需求。业务方问现在在线多少人其实不要求100%精确但他们的真实心理预期是很快看到大概趋势。rea的做法是两条路并行。一条路是实时计数。聚合服务从NATS消费事件后按照事件名分钟级时间窗口在Redis里做自增计数器。比如keyrea:count:button_click:202506131420值就是这一分钟内这个事件的数量。另一个定时任务每分钟把Redis中的值异步写入PostgreSQL的汇总表作为长期趋势的依据。这样查询端想看最近五分钟的趋势直接查Redis毫秒级返回。另一条路是明细归档。同样从NATS消费事件但这条消费者专门负责把原始事件写入PostgreSQL的rea_events表供精确查询和下钻使用。为什么要拆成两条消费链路而不是消费一条然后既聚合又入库因为两种操作的耗时差异很大。写数据库涉及磁盘IO和索引维护耗时不稳定Redis自增是纯内存操作耗时非常稳定。如果把它们混在一个管道里一次慢查询就能拖住实时计数导致报表的实时名存实亡。分而治之让实时链路尽可能短是rea一个很核心的设计决策。3.4 查询API三个缓存的配合查询API直接面向内部看板稳定性必须保证。rea的做法是三层缓存。第一层是本进程内的LRU缓存适合查同一个事件、同一个时间窗口的重复请求。第二层是Redis存的是分钟级别的预聚合数据查询时把时间段拆成分钟做聚合再在内存里合并。第三层才是数据库只有当窗口跨度过大或者需要按事件名做联合过滤时才会执行SQL查询。为什么不能只用数据库我实测过当明细表数据到了几百万行一条按事件名时间范围做COUNT的SQL在PostgreSQL里大约需要200毫秒到1秒这对于网页接口来说还可以接受。但活动大屏上的实时看板可能每5秒自动刷新一次还叠加了三四个图表如果全部打到数据库高峰期查询线程会被占满连采集服务的连接都会被拖累。加入缓存之后大部分查询落到了内存数据库压力直线下降。一个典型的API返回结构{ event: button_click, granularity: 1m, points: [ { time: 2025-06-13T14:20:00Z, count: 1200 }, { time: 2025-06-13T14:21:00Z, count: 1345 } ], queryTimeMs: 12 }3.5 权限与数据治理内部工具也有底线内部工具最容易犯的毛病是能用就行权限和安全往后放。rea虽然只对内部开放我还是花了一点时间做了几件基础的事避免以后变成大麻烦。第一事件接入按项目隔离。不同的业务项目使用不同的token事件数据里带上project_id查询时默认按当前用户可访问的项目过滤。第二userId在采集端就做哈希处理数据库里不存明文。虽然会影响精确识别单个用户的能力但对内部事件分析来说我们更关心群体行为哈希后的ID足够算留存和漏斗。第三数据保留策略。明细表只保留90天汇总表保留两年过期数据每天由定时任务清理。这条策略不是为了省存储而是为了在有人问你们凭什么存着我两年前的操作记录时我们有一个明确且合规的答案。4. 上线前后踩过的三个坑完整排查链路4.1 时间戳错乱设备本地时间把报表搞成了乱码第一个坑在试运行第二天就出现了。我打开看板发现凌晨两三点出现了一大片page_view事件而当时明显没有多少用户在访问。再看明细表某些事件的occurred_at比received_at晚了整整8小时还有一些事件的occurred_at在未来。第一反应是数据源有问题。我检查了采集服务和数据库的系统时间同步服务器时间都是准的又看了消费服务的日志发现它写入的received_at是服务端接收时间正确无误。问题只能出在事件本身的time字段。前端埋点代码用的是Date.now()也就是设备本地时间。如果用户手机的时间设置错误或者人在海外没有校准时区事件时间就会偏差很大。修复方案分两步。第一步将所有看板的时间计算基准改为received_at也就是服务端接收时间保证什么时间到达系统是可信的。第二步仍然保留occurred_at作为业务分析用的体验发生时间但是消费端增加了一个清洗逻辑如果occurred_at和received_at相差超过24小时就把occurred_at视为非法用received_at覆盖并给事件打上一个clock_skew标记。这样既保住了绝大多数准确的时间戳又不会让少数设备问题污染全局报表。这个坑给我的启示是任何带有设备端时间戳的系统都必须建立一个接收时间和发生时间的双轨制并且要有明确的校准策略。只信设备时间早晚会被雷到。4.2 高峰期写入阻塞队列共享带来的连锁问题第二个坑出现在一次推广活动流量高峰。数据库连接的等待时间瞬间拉高表象是明细表写入变慢进而导致看板上的实时指标在高峰期掉了20分钟左右的数据。我当时的排查链路是先看消费服务日志发现大量batch flush timeout再看NATS上的堆积数量超过正常水位十倍接着排查采集服务的队列积压情况发现内存队列已经满了并且Drop日志在大量刷屏说明非核心事件正在被丢弃。按理说丢弃之后积压应该缓解但指标还是掉。最后我打开了整个采集服务的线程池统计才发现问题根源不是消费端慢而是采集服务里接收线程和flush线程共用同一个线程池高峰期几条大SQL查询把线程池占满之后连事件写入NATS的线程也被阻塞了。也就是说慢查询通过线程池竞争反向拖住了整个入口链路形成连锁反应。修复措施是把IO密集的写管道和计算密集的SQL查询彻底分离。采集服务只用独立的一小簇线程处理接收事件写NATS另一簇线程专门处理查询或批量聚合。同时为NATS生产通道单独设置积压上限超过阈值时优先丢弃非核心事件而不是让所有线程一起去抢连接。这次改动之后再遇到流量高峰实时链路基本稳住了。排查线上问题一定要先画出请求的完整链路看看事件从入口到落库要经过哪几个队列、哪几个线程池所有共享资源的竞争都可能成为牵一发动全身的瓶颈。4.3 指标翻倍事件重放和缺少幂等键第三个坑是报表数据变成真实的四倍。我排查了一个下午才找到原因事后看相当典型。某天活动看板上显示按钮点击量突然涨到了平日的四倍直觉告诉我这个暴涨不合理但前端流量统计并没有明显异常。排查过程是从明细数据开始倒查的。我先查最近一小时哪个事件增长最猛很快锁定了click_feed_button然后对这一小时的事件按事件ID做去重发现ID的重复率接近75%。也就是说看板上大部分数字是重复上报带来的。继续往前看采集和消费服务都没有重复消费的逻辑消费确认机制也正常于是我打开了埋点的发送日志。真相浮出水面页面在弱网环境下fetch失败后SDK会自动重试但重试时没有重新生成事件ID每次失败重试都会把同一个事件再次发送。更糟的是在某些浏览器里页面关闭时多次触发重发逻辑一条事件能被送出去三四次。修复很简单前端SDK在创建事件的时候保留同一个id但消费端做幂等写入。数据库的rea_events表给id字段加了主键重复插入直接冲突报错汇总计数那里也按事件ID做了一次Redis的布隆过滤器重复事件不会再进入计数。后端的幂等设计看起来像是多此一举但在有网络重试机制的分布式系统里这是人命关天的事。5. 实测效果和我的最终取舍5.1 一组可复现的压测数字项目上线稳定运行一个月后我做了一组简单的压测目标是回答这套轻量级方案到底扛不扛得住。测试环境是两台2核4G的云主机一台跑采集服务和NATS一台跑PostgreSQL和聚合服务。压测工具模拟客户端以每秒100到5000的速率发送事件持续10分钟。结果汇总如下发送速率事件/秒采集服务CPU数据库CPU端到端P50延迟端到端P95延迟2005%8%400ms800ms100015%20%600ms1.2s300030%45%800ms2.8s500045%70%1.2s4.1s注意这里的端到端延迟指事件从浏览器发出到进入聚合计数的时间整体在秒级满足当初秒级延迟的设定。一旦速率超过5000数据库IO开始吃紧采集服务内存也会上涨但这已经超出rea预期的目标范围了。如果你的业务量级远高于这个数就该考虑ClickHouse和真正的流处理框架了。5.2 和商业方案和开源重方案的对比账项目做完之后我把rea和一个中等价位的商业SaaS方案做了个粗略对比。按照我们每个月大约5000万事件的体量商业方案的年费大约在五位数到六位数人民币这个量级还不包括数据出域可能带来的合规评估成本。rea的成本主要是两台云主机和一个人一个月大约30%的工作量整体低一个数量级。当然商业方案带来的价值也不能光看钱比如它有一堆现成的漏斗、留存、热力图分析不用自己开发。但那些功能对当时的我们来说是低频功能我们真正高频用的只有实时趋势、事件排行、基础漏斗三个。为了低频功能付出高频成本不划算。开源重方案那边也是一样的道理。如果当初直接上完整的流计算栈光是把环境撑起来、保证数据不丢就得花掉比业务开发更多的时间。这也引出了我的一个核心观点实时分析系统的复杂度应该跟着业务量级走而不是跟着技术潮流走。5.3 项目后续的扩展思路rea目前是能用状态但我知道它离一个完整的数据平台还差很多。这里列几个我接下来打算做的方向。第一是漏斗分析。目前只能对单个事件做计数还没有把多个事件按用户ID串联起来算转化率。我打算基于明细表写一个专门的事件序列查询接口用哈希后的用户ID做关联预计三到五天能完成。第二是异常检测。既然有了Redis里的分钟级计数完全可以写一个简单的检测器当当前分钟的事件数和过去7天同一分钟的中位数相比偏离超过三倍时推送一条告警到内部群里。这个功能对活动期实时监控特别有用不必等数据部门发现异常。第三是历史数据迁移到ClickHouse。等到明细表超过一亿行PostgreSQL的聚合查询会越来越吃力届时把rea_events同步到ClickHouse查询接口只在扫描型查询时切换数据源其余逻辑不用动。因为表结构从一开始就是按分析场景设计的这个迁移成本很低。到这个阶段我已经把rea从临时救火的脚本正式变成了一件持续迭代的工具。最后分享一点个人的实际体会。做这类内部工具最大的瓶颈不是技术而是需求边界。一开始我也想把功能做全后来发现每个新增的顺手功能都会带来额外的复杂度。rea之所以能在三周内上线并稳定运行靠的正是开始时那句边界宣言——轻量级、秒级延迟、关键事件、面向内部。如果你也想搭一套类似的实时分析系统不妨先写下这三行边界再开始选型。很多时候知道什么不做比知道做什么更值钱。
返回列表
PREV
查看更多资讯
NEXT
返回资讯列表