实时指标计算
概述
实时指标是数据价值的直接体现:PV/UV、GMV、留存、漏斗。不同指标计算难度不同,难点集中在实时去重与跨时间计算。本文给出常用实时指标的 Flink 实现方案。
一、指标分类
| 指标 | 类型 | 计算难度 |
|---|---|---|
| PV | 计数 | 低 |
| GMV | 求和 | 低 |
| UV | 去重计数 | 中 |
| 留存 | 跨天关联 | 高 |
| 漏斗 | 步骤转化 | 中 |
| 指标属性 | 说明 |
|---|---|
| 可加性 | 分区可合并(求和) |
| 半可加 | 跨维度受限(UV) |
| 不可加 | 需明细(留存) |
二、计数与求和
2.1 PV
sql
SELECT COUNT(*) AS pv FROM page_views
GROUP BY TUMBLE(biz_time, INTERVAL '1' MINUTE);2.2 GMV
sql
SELECT SUM(amount) AS gmv FROM orders
GROUP BY TUMBLE(biz_time, INTERVAL '1' MINUTE);| 要点 | 说明 |
|---|---|
| 增量聚合 | 窗口聚合状态小 |
| 事件时间 | 按业务时间 |
三、UV 实时去重
3.1 精确去重
sql
SELECT COUNT(DISTINCT user_id) AS uv FROM page_views;| 方案 | 说明 |
|---|---|
| 状态去重 | MapState 记 user_id |
| 内存集 | 适合量小 |
| 两阶段 | 局部去重 + 合并 |
精确去重:状态存所有 user_id
量大时状态大 → 内存/恢复压力3.2 近似去重
| 算法 | 说明 |
|---|---|
| BloomFilter | 判断存在,有误判 |
| HyperLogLog | 基数估计,误差约 1% |
Flink 内置 HLL 近似 UV:
误差可控、状态小3.3 选型
| 场景 | 方案 |
|---|---|
| UV 小 | COUNT DISTINCT |
| UV 大、可近似 | HLL |
| 精确必须 | 状态去重 + 明细落湖 |
3.4 分桶去重
大用户量:
按 user_id 分桶 → 各桶去重 → 汇总
并行分散状态四、留存计算
4.1 定义
次日留存 = 今天活跃且明天活跃的用户 / 今天活跃用户4.2 实现
sql
-- 活跃用户表(按日期去重)
CREATE VIEW act AS
SELECT user_id, DATE_FORMAT(biz_time, 'yyyy-MM-dd') AS dt
FROM page_views GROUP BY user_id, dt;
-- 留存:T 日活跃用户是否 T+1 活跃
SELECT a.dt AS base_dt,
COUNT(DISTINCT CASE WHEN b.dt IS NOT NULL
THEN a.user_id END) AS retain_users,
COUNT(DISTINCT a.user_id) AS base_users
FROM act a
LEFT JOIN act b
ON a.user_id = b.user_id AND b.dt = DATE_ADD(a.dt, 1);4.3 留存难点
| 难点 | 处理 |
|---|---|
| 跨天 | 维护每日活跃集合 |
| 状态大 | 按天分区状态 |
| 指标滚动 | 窗口式留存 |
留存状态:
保留 N 天活跃集合
每日更新,过期清理五、漏斗计算
5.1 定义
漏斗:按步骤统计转化率
访问 → 加购 → 下单 → 支付5.2 实现
| 方式 | 说明 |
|---|---|
| 多流 Join | 各步骤事件关联 |
| 状态标记 | 记录用户最大完成步骤 |
| CEP | 事件序列匹配 |
5.3 状态标记法
按用户记录当前进度:
访问 → progress=1
加购(progress≥1)→ progress=2
下单(progress≥2)→ progress=3
支付(progress≥3)→ progress=4
各步骤计数 = progress ≥ 步骤的用户数| 步骤 | 计数逻辑 |
|---|---|
| 访问 | progress >= 1 |
| 加购 | progress >= 2 |
| 下单 | progress >= 3 |
| 支付 | progress >= 4 |
5.4 CEP 方式
模式:访问 → 加购 → 下单 → 支付(事件时间窗口)
匹配到完整序列 → 计入漏斗六、实时指标架构
6.1 计算引擎
Kafka 明细 → Flink 实时计算 → 指标输出| 输出 | 说明 |
|---|---|
| Kafka | 继续流转 |
| Doris/ClickHouse | 查询/大屏 |
| Redis | 实时值 |
6.2 与离线对账
实时指标 vs 离线指标:
趋势一致、允许小误差
对账修正系统性偏差七、指标一致性
| 问题 | 处理 |
|---|---|
| 口径不一 | 统一指标定义 |
| 重复 | 幂等/去重 |
| 迟到 | 水位线 + 修正 |
| 偏差 | 与离线对账 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| UV 偏大 | 重复数据未去重 |
| 留存率异常 | 活跃口径不一致 |
| 漏斗断链 | 步骤事件缺失 |
| 去重状态大 | 用 HLL 或分桶 |
| 指标漂移 | 与离线对账修正 |