ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

实时数仓长周期去重指标实践:Flink + Doris Bitmap 全局字典方案

实时数仓长周期去重指标实践:Flink + Doris Bitmap 全局字典方案 如果你参与过实时数仓建设大概率被这类需求追过运营要“近30天活跃会员数”风控要“近90天下单用户数”数据分析师隔三差五要一份“历史累计去重访客”。这些指标有一个共同的名字——长周期去重指标也是实时数仓里公认难啃的骨头。难在哪它不是把离线的count(distinct)平移过来就完事实时流上做长周期精确去重既要压住状态膨胀又要扛住口径校验很多团队在这个环节翻过车。这篇我把自己从零搭建长周期去重指标体系的完整过程拆开讲包括指标口径分析、方案选型对比、全局字典构建、滚动窗口落地、上线调优和效果复盘。适合已经在做实时数仓、被 UV 类去重指标折磨过的开发同学参考如果你正准备从 0 起步建设实时数仓这篇也能帮你少走不少弯路。1. 先想清楚口径长周期去重指标到底难在哪1.1 去重指标不是只有“用户数”一种拿到需求先别急着写代码第一件事是确认“去重主体”和“统计范围”。最常见的去重主体是用户 ID但实际业务里还有设备 ID、订单 ID、支付流水号、渠道 ID 等。不同主体对应完全不同的数据链路和状态设计。举个例子一个典型的电商实时大屏“今日下单用户数”只按用户 ID 去重“今日订单数”按订单 ID 去重“今日新客数”还要和“历史是否下过单”做碰撞。前两个看似简单实际上数据迟到、重复上报、多端登录都会导致口径偏移。更麻烦的是“长周期”一个用户过去 30 天每天都活跃那他在“近30天UV”里只能算一次而不是 30 次。这个逻辑听起来理所当然但真正在流上实现时会发现你必须在长达 30 天的时间窗口里记住“这个用户来过”。1.2 固定周期、滚动周期和历史累计三种口径差别巨大长周期去重指标按时间窗切常见有三种自然周期截止型比如“本月至今 UV”时间窗口是 1 号到当前时刻每天计算一次窗口固定。滚动窗口型比如“近 30 天 UV”每个自然日算过去 30 天窗口每天往前滑一天。这是最常被业务点名、也最难实时化的口径。历史累计型比如“累计注册用户数”“上线至今去重用户”窗口相当于从业务启动到当前且只增不减。很多人一开始只关注“去重”忽略了“窗口形态”结果方案设计完发现滚动窗口做不了。固定周期可以每天清零重算历史累计可以只维护一个全局状态唯独滚动窗口要求你把“过去 N 天”这个区间内所有出现过的用户都记住而且 N 越大要记的东西越多——这就是长周期去重最核心的矛盾。1.3 容易被忽略的口径细节即使定好了“近 30 天 UV”颗粒度也没完自然日还是自然周按事件发生时间还是按数据到达时间跨凌晨 0 点的一条活跃记录算昨天还是今天同一个用户同时在 App 和 H5 产生行为算一个还是两个有些业务还要“先按 device_id 去重再按 user_id 去重”。这些细节在离线数仓里通常靠 SQL 的where event_date between ...一次性算完但在实时链路里一旦有边界条件没锁死后面所有下游指标全部对不上。我踩过最典型的坑是时区埋点事件时间戳用的 UTC业务看板用的是东八区实时统计和离线 T1 对账始终差几个小时最后排查半天发现是时区偏移没处理。1.4 状态爆炸的根因拆解我们做实时数仓去重说白了就是在流上维护一组“已经见过的 key”。一天的去重可能只要存几千万个 key30 天去重就要存 30 天里所有出现过的 key——注意这里不是简单乘 30因为活跃用户的重合度很高但又不能赌“重合度一定高”而放弃精确存储。用经典的 Flinkcount(distinct user_id)方案时状态需要保存窗口内全部 user_id 明细key 的量级随窗口长度线性增长内存和磁盘都扛不住。更麻烦的是RocksDB 的 key 一多checkpoint 经常超时恢复也慢故障时很难在 SLA 内追上进度。所以长周期去重指标建设的本质不是“怎么去重”而是“怎么用可控的存储换取可接受的查询性能”。2. 方案选型四条路我都试过最后定在 Flink Doris Bitmap2.1 Flink SQL 状态去重最简单但容器装不下刚接需求时我第一反应是用 Flink SQL 直接做SELECT COUNT(DISTINCT user_id) AS uv_30d FROM dwd_user_behavior WHERE event_time NOW() - INTERVAL 30 DAY;语法上没毛病跑起来很快就出问题。Flink 内部要给每个算子维护状态COUNT(DISTINCT ...)会把所有 user_id 都放进状态里。业务量级小还能忍一旦每天 DAU 上千万30 天窗口里要去重的用户数轻松过亿。RocksDB 单机状态几十 GBcheckpoint 频繁超时背压一路传到 Kafka 消费端实时链路变成“实时延迟”。我当时也试过给状态配 TTL比如STATE_TTL设 30 天。TTL 确实能回收过期 key但问题是 TTL 是按“状态写入时刻”算的一个用户第 29 天活跃了它的状态又被续到第 59 天实际要占用近 60 天窗口量的状态。窗口越大这种“隐性膨胀”越严重。结论Flink 状态只适合小周期或低基数去重不适合长周期大基数。2.2 Redis Set成本高恢复和大 key 让人崩溃第二个方案在团队内部呼声很高用 Redis Set每天一个 keySADD往里塞 user_id查询时SUNION合并再SCARD算基数。思路直接但问题同样明显。首先是成本。一个 user_id 假设 16 字节存 1 亿个用户就是 1.6GB 起步SD 内存压力巨大。其次是稳定性单日 key 过千万后SINTER/SUNION的耗时开始不可控大 key 在 AOF 重写和主从同步时容易阻塞。更重要的是数据恢复Redis 故障后要重新灌 30 天的数据回放才能恢复指标这个时间窗口业务完全不可接受。当时我们评估了一下服务器预算直接放弃。Redis 适合做缓存、做轻量级计数不太适合做长周期海量去重的主体存储。2.3 离线/实时混合口径容易对不上还有一种方案是“实时只算今日历史走离线聚合最后拼装”。比如今天凌晨用离线任务算好“截至昨天的 30 天 UV”实时环节只加上今天的增量。这个方案在早期确实能用但业务一旦要看“任意一天往前推 30 天”的指标离线就得把每个历史日期的 30 天窗口全算一遍存储和调度成本急剧上升。而且“实时增量”和“离线历史”两份口径拼接时经常出现重复计数或漏计两边对不齐时定位原因非常痛苦。2.4 最终选择全局字典 Doris Bitmap和团队反复评估后我们确定了技术路线Flink 负责实时 ETL 与字典关联Hive 负责离线构建全局字典Doris 负责存储 Bitmap 和对外提供查询服务。选 Doris 有几个原因原生支持Bitmap类型和bitmap_union聚合建表时能声明聚合模型写入即按 key 聚合对分钟级实时写入非常友好查询端 SQL 语义简单业务侧直接count(bitmap_union(...))就能拿到精确去重值。和 ClickHouse 相比Doris 在精确去重的关联查询、多表 join 和 Update 模型上更贴合我们数仓的使用习惯。下面是当时方案选型时做的对比表我直接贴出来方案存储主体查询性能状态/空间占用成本适用场景Flink State原始 user_id 明细一般随窗口线性膨胀中短周期、低基数Redis Set原始 user_id 明细大 key 后劣化高高小规模临时需求离线实时混合离线表实时状态一般中中口径要求不严格Doris Bitmap整数 ID 位图毫秒到秒级压缩后极小可控长周期、大规模精确去重3. 全局字典是怎么建起来的这是整个方案的地基3.1 Bitmap 为什么非要整数映射Bitmap 的原理说起来很简单把“是否出现过”映射到一个个二进制位第 N 位为 1 表示 ID 为 N 的元素出现过。它只能高效处理整数 ID而业务传入的 user_id 是 UUID 或长字符串直接用字符串做位图没有任何意义。所以必须先建一张“全局字典”把每一个字符串 user_id 映射成唯一的整数 IDuser_id - guid。这张字典的完整性和唯一性直接决定了 Bitmap 方案的正确性。字典漏了一个用户去重数就少一个字典里同一个用户映射了两个 ID去重数就虚高。长周期去重指标的误差大概率不是在引擎侧而是在字典侧。3.2 离线全量字典Hive 里用自增序号生成字典的底表我选在 Hive 构建。核心逻辑是把所有历史 user_id 拉全去重后按序分配整数 ID-- dim_user_guid 为字典表uid 为原始字符串guid 为唯一整型 CREATE TABLE dim_user_guid ( uid STRING, guid BIGINT ); -- 初始构建对全量用户去重后分配自增 ID INSERT OVERWRITE TABLE dim_user_guid SELECT uid, row_number() OVER (ORDER BY uid) - 1 AS guid FROM ( SELECT DISTINCT uid FROM ods_user_all ) t;但生产环境不是一次性构建每天都有新用户产生。增量更新时要注意不能重复分配也不能重新全量覆盖因为线上 Bitmap 里已经引用了历史 guid重建字典会导致历史数据全部作废。增量逻辑这样写-- 假设已存在 dim_user_guid INSERT INTO TABLE dim_user_guid SELECT new_uid.uid, old.max_guid row_number() OVER (ORDER BY new_uid.uid) AS guid FROM ( SELECT uid FROM ods_user_all WHERE uid NOT IN (SELECT uid FROM dim_user_guid) ) new_uid CROSS JOIN ( SELECT COALESCE(MAX(guid), -1) AS max_guid FROM dim_user_guid ) old;离线字典每天凌晨调度一次产出的结果同步到 Doris 维表或 HBase供实时任务查询。3.3 实时新增 ID 的分配缓存 批量回填离线字典只能覆盖“昨天之前”的历史用户当天新产生的 user_id 必须由实时链路即时分配 guid。我们在 Flink 里做了一个 RichFlatMap 算子内部维护一个 LRU 缓存流程是收到一条 (stat_date, user_id)先在本地缓存查 uid → guid缓存 miss则查外部字典存储我们用 HBase 存储 uid 到 guid 的映射HBase 里查到了写回本地缓存直接使用HBase 里查不到说明是新用户向 Redis 请求一个自增 IDINCR同时写回 HBase 和本地缓存。这里有个并发风险必须处理同一个新 user_id 如果同时被两个并发子任务处理可能分配出两个不同的 guid。我们当时在 HBase 写入时做了“先查再写”的防护同一 uid 若已存在则丢弃新分配的 ID另一种方案是启用 Redis 分布式锁但为了性能实际我们更倾向用 HBase 的 check-and-set 原子操作来兜底。这个细节没处理好的话去重数会虚高而且很难排查。3.4 字典不一致会带来哪些隐形 Bug字典层面踩过的坑我印象最深不夸张地说它比引擎本身的 Bug 更难发现。一个问题是字典分裂离线重建字典当天实时链路还在用旧字典两边各分配各的 guid导致同一 uid 存在两条映射记录。表现就是当天实时 UV 和离线 UV 对不上且差值不固定。解决方法是把字典的离线更新和实时分配串成一条线夜间离线任务结束后先把最新字典同步到 HBase再让实时任务开始消费新增用户过渡期内实时任务统一走 HBase 查询。另一个问题是字典只增不删业务里会有销户、合并账号的场景。为了 Bitmap 的稳定性我们约定字典表不物理删除任何记录只在业务侧打标。因为历史 Bitmap 里已经引用了这些 uid 对应的 guid删了之后不仅无法回溯明细还会破坏所有长周期指标的连续性。4. 滚动窗口的落地按“天”切 Bitmap用“并集”代替“大状态”4.1 核心思路把窗口换成多日 Bitmap 的并集字典建好之后最关键的架构决策来了长周期滚动窗口不能直接用一个超大的 Bitmap 去维护而是按自然日切分成单日 Bitmap查询时再把多日 Bitmap 做并集OR 操作。举个例子“近 30 天 UV”就是今天和前 29 天的每日 Bitmap 做按位或得到一个 30 天合并 Bitmap再数里面有多少个 1。这样做的好处是存储从“30 天窗口的重复状态”变成“30 个单日 Bitmap”每天的 Bitmap 是天然增量凌晨 0 点自然切换完全不需要 TTL 清理逻辑。单日 Bitmap 的生成也符合数仓分层习惯DWD 层是明细DWS 层按天聚合。每天只保留一个日粒度 Bitmap过期的日分区可以定期归档或淘汰状态没有跨天包袱。4.2 DWS 层表设计Doris 聚合模型Doris 建表的关键点是BITMAP_UNION聚合类型。以渠道维度为例DWS 层每日去重表长这样CREATE TABLE dws_uv_daily ( stat_date DATE, channel_id VARCHAR(32), user_bitmap BITMAP BITMAP_UNION ) AGGREGATE KEY(stat_date, channel_id) DISTRIBUTED BY HASH(channel_id) BUCKETS 12 PROPERTIES ( replication_num 3 );Flink 侧在写入 Doris 之前已经通过全局字典把 user_id 转换成整数 guid导出时按 (stat_date, channel_id, guid) 输出。Doris 在导入过程中会按照聚合键自动把同一分区的 guid 集合bitmap_union成一个 Bitmap。也就是说我们不需要在 Flink 里自己攒 Bitmap只要把“日 维度 整数 ID”的中继明细写进去聚合成 Bitmap 的工作完全交给 Doris。4.3 近 N 天 UV 的查询与预聚合有了每日 Bitmap滚动窗口查询就是一条标准 SQLSELECT channel_id, bitmap_count(bitmap_union(user_bitmap)) AS uv_30d FROM dws_uv_daily WHERE stat_date BETWEEN 2024-11-01 AND 2024-11-30 GROUP BY channel_id;这里bitmap_union在 Doris 里就是做多行 Bitmap 的按位或聚合等价于把 30 天的每日 Bitmap 合并成一个大 Bitmap再bitmap_count统计基数。实测千万级日活下30 天窗口的查询能控制在 2 秒内比 Flink 状态方案快了不止一个量级。如果业务方高频查询的窗口比较固定比如 7 天、30 天、90 天我建议再建一层预聚合表用定时任务把对应窗口的 Bitmap 合并结果写进去。查询直接读预聚合结果把计算压力从查询侧挪到离线调度侧看板和即时分析都吃得住。4.4 离线对账是必须做的实时链路再自信对不上离线账就没有说服力。我们每天跑两个对账任务自然日对账Hive 当天count(distinct user_id)对比 Doris 当日bitmap_count(...)滚动窗口对账Hive 对近 30 天明细去重对比 Doris 的 30 天预聚合结果。对账发现差异时先查时区、再查字典、最后查迟到数据。我这边实际跑下来字典稳定后自然日对账能精确到 0 误差滚动窗口对账偶尔有几十到几百的偏差基本都来自埋点数据迟到属于可接受范围。5. 上线后我踩过的六个坑5.1 Checkpoint 超时状态拆分比调参更重要方案刚上线时Flink 作业还是背着不少状态因为字典关联算子里缓存了 uid → guid 映射。checkpoint 经常跑到一半就超时一度怀疑是 RocksDB 的问题。后来定位到根因是缓存没有做序列化控制LRU 缓存里的映射表被默认纳入算子状态数据量一大snapshot 就慢。解决方式是把“本地缓存”改为transient不参与 checkpointHBase 才是唯一可靠数据源本地缓存只是加速。改造后 checkpoint 大小从 20GB 降到几百 MB稳定性立竿见影。这给我们的经验是能放进外部存储的状态就不要留在 Flink 里实时任务的状态只放“过程态”不要放“结果态”。5.2 位图膨胀RoaringBitmap 也怕垃圾数据Bitmap 的压缩率依赖 ID 分布的连续性。某天我们导入时误把字符串 user_id 直接塞进了 Bitmap 列没有先做字典映射结果 Bitmap 不但没压缩反而比原始字符串集合还大。后来查了 Doris 的导入日志才确认字段映射错了本该用user_bitmapto_bitmap(guid)的地方写成了user_bitmapto_bitmap(user_id)。这个错误会造成整个调度集群的写入放大。排查方式是用 Doris 自带的 profile 看单个 tablet 的大小一旦出现异常膨胀优先检查映射函数。5.3 字典孤儿 ID 导致历史数据不可回溯这是字典层面最麻烦的事。某次表重构时我们把 Hive 字典表OVERWRITE重建了guid 分配顺序发生变化历史 Bitmap 里所有“用户”的编号全部失效。当日实时指标没感觉等业务要查“30 天趋势”时发现前 29 天全部错乱且无法通过字典反查明细。从此我们定了一条铁律字典表只允许 append不允许 overwrite任何导致 guid 语义变化的操作都必须保留完整历史映射表。Bitmap 从设计上就不支持“按值反查”一旦映射断裂数据等于永久损坏只能从源头避免。5.4 Stream Load 小文件拖垮导入性能Flink 写入 Doris 时的批次设置也是容易翻车的点。最开始我图传输实时Sink 每 5 秒就 flush 一次Doris BE 不断创建新版本后台 compaction 永远跑不过写入速度查询时经常要合并几十个版本性能直线下降。调整方案是把 Stream Load 的批次调大每次至少攒 64MB 或 60 秒再 flush。导入频率降下来之后BE 的 compaction 压力明显缓解。写 Doris 这类实时 OLAP最忌讳“高频率小批次”宁可稍等几秒再批量写也不要制造一堆小版本。5.5 分区分桶没规划好查询照样慢DWS 表的分区按自然日这是没问题的。但分桶最开始我按stat_date哈希结果所有数据都挤到少数几个桶里数据倾斜严重查询并行度上不去。后来改成按channel_id分桶并把桶数从 10 调整到 24查询耗时稳定在秒级。分桶键的选择要考虑“查询过滤条件”和“聚合粒度”如果业务主要按渠道查询分桶键就应该带渠道字段如果纯看总体 UV其实单分区直接全扫描也能接受但分桶要保证均匀。5.6 预聚合表刷数踩的调度坑预聚合表跑批调度时我们一开始每天只算“今天往前推 30 天”的结果忘了处理“昨天出过数据、今天要修正”的场景。结果某天上游出现迟到数据补偿预聚合表却用了覆盖写把已经正确的历史窗口结果又刷成了不完整版本。修正方法是把预聚合调度做成“全窗口重算”每次调度都重算最近 31 天的所有日分区再只更新结果表里涉及变化的窗口分区。数据量不大但正确性彻底兜住了。6. 效果复盘与这套方案的适用边界6.1 我这边跑出来的实际数据以我们当时的业务量级举例日活跃用户约 3000 万30 天累计去重用户约 1.8 亿。方案上线稳定后几个关键指标是这样的Flink 作业状态从之前方案的 20GB 级降到不足 1GBcheckpoint 大小缩减 95% 以上Doris 每日新增一个日粒度 Bitmap单日 Bitmap 存储约 20~80MB30 天全部 Bitmap 占用控制在 2GB 内30 天窗口查询 P95 耗时 1.8 秒7 天窗口不到 1 秒自然日 UV 与离线 Hive 对账零误差30 天滚动窗口误差率低于万分之一。这个效果的根因在于Flink 只负责“把字符串映射成整数”和“轻量 ETL”真正吃存储和计算能力的 Bitmap 聚合下沉到了 Doris。实时计算的压力被极大释放数据服务能力反而更强了。6.2 这套方案不适合哪些场景再好的方案也有边界。Bitmap 全局字典这套路在以下场景会很难受超高基数且不可映射如果去重主体本身是几十亿量级的随机设备指纹全局字典膨胀到比原始数据还大字典的构建和查询就变成新的瓶颈需要频繁反查明细Bitmap 只适合算基数业务如果要求“给我这 30 天活跃用户的手机号列表”你必须把 guid 反查回 uid 再关联业务库这条链路比直接存明细更绕窗口极短如果只是“近 1 小时 UV”Flink 状态去重反而更直接没必要为了引入 Bitmap 增加字典维护成本。6.3 后续还能怎么做这个架构跑顺之后扩展方向很清晰累计去重指标可以做成“历史至今 Bitmap 每日增量 Bitmap 合并”留存指标可以复用每日 Bitmap 做集合运算比如“30 天活跃用户里有多少今天也活跃”甚至跨主题去重也能通过多个 Bitmap 的 AND/OR 组合实现SQL 层面基本不用大改。我个人的习惯是每上线一套新指标先把“口径 字典 对账”三个文档同步给下游再放开查询权限。事实证明这个顺序能挡掉大量“数据是不是错了”的咨询。最后分享一个细节全局字典表一定要留历史版本和归档策略别图省事只保留最新一份。长周期数据回溯时它就是救命的底牌。这套方案的后续演进等我把累计去重和留存指标做完再写一篇展开聊。
返回列表