StarRocks retention 聚合函数详解:用户留存率计算从 SQL 语法到底层位压缩实现

StarRocks retention 聚合函数详解:用户留存率计算从 SQL 语法到底层位压缩实现 StarRocks retention 聚合函数详解用户留存率计算从 SQL 语法到底层位压缩实现【免费下载链接】starrocksThe worlds fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks本篇指南围绕 StarRocks 的retention聚合函数展开先完整讲解其语法、参数、返回值规则与留存率计算的完整实操案例再结合 BE 端 C 实现剖析最多 31 个条件背后的 64 位位压缩状态机帮助读者既会写留存分析 SQL也能理解该函数在分布式聚合中的真实执行机制。retention 函数概述retention是 StarRocks 提供的聚合函数aggregate function用于计算指定时间段内的用户留存率。它的独特之处在于接受一个条件数组1 到 31 个布尔条件对分组内每一行逐条求值最终返回一个由 0 和 1 组成的数组供后续按元素下标统计满足某条件链的分组数进而算出留存率。典型应用场景次日/次日以后留存率第 1 个条件通常是首日行为基线后续条件表示第 N 日行为用sum(r[N])/sum(r[1])得到相对首日的留存比例用户转化漏斗式留存分析例如浏览商品页→下单两个行为是否同时发生与GROUP BY user_id配合按用户维度聚合行为后再汇总无需自连接self join多个日期子查询。在官方文档中retention归属于聚合函数aggregate functions章节文档原文位于 retention.md。语法、参数与返回值语法ARRAY retention(ARRAY input)参数input条件数组最多传入31 个条件多个条件之间用逗号分隔。从前端解析代码可以确认参数必须是ARRAYBOOLEAN类型否则分析阶段直接报错见下文FE 类型校验一节。返回值规则返回一个由 0 和 1 组成的数组元素个数与输入条件数相同求值从第一个条件开始条件求值为 true 返回 1否则返回 0如果第一个条件不为 true则当前位置及其后所有位置一律置 0——即后续元素隐含以第一个条件成立为前提的语义。这一规则在 BE 源码的finalize_to_array_column中有对应实现后文详述也就是说结果数组第 N 个元素实际表示条件 1 成立 且 条件 N 成立。完整实操案例以下案例完整继承官方文档示例可直接复制执行。步骤 1建表并插入数据CREATE TABLE test( id TINYINT, action STRING, time DATETIME ) ENGINEolap DUPLICATE KEY(id) DISTRIBUTED BY HASH(id); INSERT INTO test VALUES (1,pv,2022-01-01 08:00:05), (2,pv,2022-01-01 10:20:08), (1,buy,2022-01-02 15:30:10), (2,pv,2022-01-02 17:30:05), (3,buy,2022-01-01 05:30:09), (3,buy,22022-01-02 08:10:15), (4,pv,2022-01-02 21:09:15), (5,pv,2022-01-01 22:10:53), (5,pv,2022-01-02 19:10:52), (5,buy,2022-01-02 20:00:50);步骤 2查询数据select * from test order by id;----------------------------------- | id | action | time | ----------------------------------- | 1 | pv | 2022-01-01 08:00:05 | | 1 | buy | 2022-01-02 15:30:10 | | 2 | pv | 2022-01-01 10:20:08 | | 2 | pv | 2022-01-02 17:30:05 | | 3 | buy | 2022-01-01 05:30:09 | | 3 | buy | 2022-01-02 08:10:15 | | 4 | pv | 2022-01-02 21:09:15 | | 5 | pv | 2022-01-01 22:10:53 | | 5 | pv | 2022-01-02 19:10:52 | | 5 | buy | 2022-01-02 20:00:50 | ----------------------------------- 10 rows in set (0.01 sec)步骤 3用 retention 计算用户留存示例 1评估用户在 2022-01-01 浏览商品页actionpv且在 2022-01-02 下单actionbuy两个条件。select id, retention([actionpv and to_date(time)2022-01-01, actionbuy and to_date(time)2022-01-02]) as retention from test group by id order by id;----------------- | id | retention | ----------------- | 1 | [1,1] | | 2 | [1,0] | | 3 | [0,0] | | 4 | [0,0] | | 5 | [1,1] | ----------------- 5 rows in set (0.01 sec)结果解读用户 1 和 5 满足两个条件返回[1,1]用户 2 不满足第二个条件返回[1,0]用户 3 满足第二个条件但不满足第一个条件返回[0,0]体现首条件不成立则后续全 0的规则用户 4 两个条件都不满足返回[0,0]。示例 2计算2022-01-01 浏览商品页的用户中2022-01-02 下单的占比即次日留存率。select sum(r[1]),sum(r[2])/sum(r[1]) from (select id, retention([actionpv and to_date(time)2022-01-01, actionbuy and to_date(time)2022-01-02]) as r from test group by id order by id) t;-------------------------------------- | sum(r[1]) | (sum(r[2])) / (sum(r[1])) | -------------------------------------- | 3 | 0.6666666666666666 | -------------------------------------- 1 row in set (0.02 sec)返回的比值即为 2022-01-02 的用户留存率sum(r[1]) 3表示首日浏览的基线用户数用户 1、2、5sum(r[2]) 2表示其中在次日下单的人数用户 1、5比值约为 0.667。FE 侧函数注册与参数校验retention在 FE 的内置函数目录中按入参ARRAYBOOLEAN→ 返回ARRAYBOOLEAN中间聚合状态为BIGINT注册见 FunctionSet.java// Retention addBuiltin(AggregateFunction.createBuiltin(RETENTION, Lists.newArrayList(ArrayType.ARRAY_BOOLEAN), ArrayType.ARRAY_BOOLEAN, IntegerType.BIGINT, false, false, false));中间的IntegerType.BIGINT参数值得注意它声明了该聚合函数在分阶段执行本地聚合 → 序列化 → 合并时的中间状态类型是一个 64 位整数。这与 BE 端用一个uint64_t压缩整个分组状态的实现一一对应意味着无论条件有多少个序列化传输的状态都只有一个整数开销恒定。参数校验在 FunctionAnalyzer.java 中完成if (fnName.equals(FunctionSet.RETENTION)) { if (!arg.getType().isArrayType()) { throw new SemanticException(retention only support ArrayBOOLEAN, arg.getPos()); } ArrayType type (ArrayType) arg.getType(); if (!type.getItemType().isNull() !type.getItemType().isBoolean()) { throw new SemanticException(retention only support ArrayBOOLEAN, arg.getPos()); } // For ArrayBOOLEAN that have different size, we just extend result array to Compatible with it }由此可以确认两条使用约束其一参数必须是数组且元素类型为布尔值写retention(actionpv, actionbuy)这种散参写法会直接报 SemanticException其二从注释对于大小不同的ArrayBOOLEAN将扩展结果数组以兼容来看可以推断当同组内各行传入的条件数组长度不一致时结果数组长度以较长者为准做兼容处理。FE 计划测试 AggregateTest.java 中也覆盖了retention([true,true])、retention([])等计划的生成用例。BE 侧实现一个 64 位整数压缩 31 个条件retention的 BE 实现位于 retention.h并在 BE 聚合工厂中以数组映射形式注册入参数组、出参数组见 aggregate_resolver_others.cppadd_array_mappingTYPE_ARRAY, TYPE_ARRAY(retention);状态结构31 个条件位 5 个长度位核心是RetentionState结构体只用一个uint64_t boolean_value表示一个分组的完整聚合状态见 retention.h// We use top 31 bits of boolean_value to indicate which condition is true; // We use the last 5 bits of boolean_value to indicate size of conditions(so 31 conditions at most). uint64_t boolean_value; // Mask is used to identify top 31 bits. static inline uint64_t bool_values[] {1UL 63, 1UL 62, 1UL 61, /* ... */ 1UL 33}; static constexpr int MAX_CONDITION_SIZE_BIT 5; // Mask is used to identify the last 5 bits. static constexpr int MAX_CONDITION_SIZE (1 MAX_CONDITION_SIZE_BIT) - 1;位布局设计如下位区间含义第 63 位 ~ 第 33 位共 31 位第 1 ~ 第 31 个条件是否曾为 true每一位对应bool_values[]中的一个掩码第 4 位 ~ 第 0 位共 5 位条件数组的长度最多 31 2^5 - 1这也从实现层面解释了最多 31 个条件这一上限的来源5 个长度位最多表示 31与条件位数量恰好对齐。udpate中还有显式保护if (array_size MAX_CONDITION_SIZE) { array_size MAX_CONDITION_SIZE; }防止越界。逐行更新条件数组按位 ORupdate是逐行chunk 内第row_num行更新状态的入口。每行传入的是已求值的布尔数组例如每行对应[pv 且 01-01 是否成立, buy 且 01-02 是否成立]RetentionState::udpate遍历该数组元素for (size_t i 0; i array_size; i) { auto ele_offset offset i; if (!null_column.is_null(ele_offset) imm_data[ele_offset]) { // Set right bit for condition. (*value_ptr) | RetentionState::bool_values[i]; } } (*value_ptr) | array_size;从这段实现可以确认其聚合语义同一个 GROUP BY 组内各行的条件结果是按位或OR累积的。也就是说只要用户 5 的多行数据中有一行满足01-02 buy该组状态的第 2 个条件位就会被置 1——这正对应示例 1 中多行INSERT数据最终归并为每用户一个结果数组的行为。同时条件元素为 NULL 时不会置位等价于该位置为 false。分布式合并与序列化只传一个 BIGINTStarRocks 聚合查询通常是本地聚合 → 跨节点交换部分聚合 → 全局合并的流程RetentionAggregateFunction对每个环节都有对应方法见 retention.hvoid merge(...) const override { const auto* input_column down_castconst Int64Column*(column); this-data(state).boolean_value | input_column-immutable_data()[row_num]; } void serialize_to_column(...) const override { down_castInt64Column*(to)-append(this-data(state).boolean_value); }serialize把 64 位状态直接作为一个Int64写出与 FE 注册的 BIGINT 中间状态类型一致merge收到远端发来的部分聚合状态后直接做位或|——由于位与位之间互不干扰条件位与长度位互不重叠长度位对同一函数恒为各分组条件数的或而各分组条件数相同位或操作即完成了两个分组的正确合并网络传输代价因此与条件个数无关始终是一个 8 字节整数。收尾首条件为 0 则后续全 0 在这里实现文档中第一个条件不为 true 时当前及其后所有位置置 0的规则在finalize_to_array_column中实现void finalize_to_array_column(ArrayColumn* array_column) const { auto size (boolean_value MAX_CONDITION_SIZE); DatumArray array; if (size 0) { array.reserve(size); auto first_condition ((boolean_value bool_values[0]) 0); array.emplace_back((uint8_t)first_condition); if (first_condition) { for (int i 1; i size; i) { array.emplace_back((uint8_t)((boolean_value bool_values[i]) 0)); } } else { for (int i 1; i size; i) { array.emplace_back((uint8_t)0); } } } array_column-append_datum(array); }最终输出时才读取bool_values[0]第 1 个条件位若其为 0直接把第 2 项起全部填 0否则逐项输出各条件位。注意这是输出时裁剪而非求值时裁剪状态位始终忠实记录各条件的真实命中情况规则只在最终物化为数组时生效——这解释了为什么文档把该规则放在Return value而非Parameters一节。使用要点与适用前提小结综合文档与源码证据使用retention时的关键要点必须配合 GROUP BY 使用它是聚合函数条件数组在分组内按位 OR 累积留存率统计的标准写法是先按用户GROUP BY出每用户的结果数组再对外层做sum(r[n])/sum(r[1])汇总条件数量上限 31由 5 位长度字段决定超出会被截断处理做 7 日留存基线 7 个后续日等场景远在限额内基线条件必须放第一个由于后续元素隐含条件 1 成立的前提若把想统计的目标行为放在第 1 位则所有后续位都可能被清零留存率将失真参数必须是布尔条件数组FE 对非ARRAYBOOLEAN入参报retention only support ArrayBOOLEAN日期比较建议用to_date(time)yyyy-MM-dd如示例所示而非直接比较 DATETIME避免时分秒干扰适用前提示例基于当前仓库文档与 OLAP 内表ENGINEolap演示行为对数据湖外表查询同样适用因为该函数是 BE 通用聚合函数与表类型无关。相关地StarRocks 还提供了同为行为序列分析的window_funnel聚合函数同样接收ARRAYBOOLEAN见 FunctionSet.java适用于带时间窗口的漏斗统计若只需要是否达成行为链 按位留存比例retention是更轻量的选择。参考文件函数文档retention.mdBE 实现retention.h、注册处 aggregate_resolver_others.cppFE 注册与校验FunctionSet.java、FunctionAnalyzer.javaFE 计划测试AggregateTest.java【免费下载链接】starrocksThe worlds fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考