算子语义规格
July 26, 2026 · View on GitHub
本文是规范性文档(normative)。 引擎实现、SDK、回归测试都以本文为准。当实现与本文冲突时,以本文为准并修正实现。
本文同时是 1.x → 2.0 的语义对照表。凡与 1.x 行为不同之处,均在对应条目下的「与 1.x 的差异」中标明,这些差异构成迁移时的不兼容点。
阅读约定
- 窗口内指当前窗口尚未过期的全部事件
- key 指由
groupbykeys求值得到的分组键;groupbykeys为空时使用全局键 - 除非特别说明,所有算子遇到
null输入一律跳过该条输入,而不是当作 0 或空串 - 窗口内无任何输入时的返回值,在每个算子的「空窗口」中单独规定
一、计算图节点
变量按 type 映射为计算图中的一类节点。一条事件进入后沿图自根向下传播,节点返回"不通过"即剪枝,该分支的下游不再计算。
1.1 event —— 事件解包
计算图的根节点,每张图有且只有一个。
- 输入:一条原始事件
- 输出:事件上下文,包含事件的全部属性,以及
app、name、key、timestamp、value五个固定字段 - 状态:无
- 恒通过:不做任何过滤
1.2 filter —— 过滤与字段派生
- 输入:上游上下文
- 输出:不通过时剪枝;通过时把 mapping 产出的派生字段写入上下文,并透传
value与timestamp - 状态:无
求值顺序是一个必须精确保留的语义:
- 当过滤条件引用的是上游已有字段时,先判条件、后做 mapping(条件不满足则跳过 mapping 的计算开销)
- 当过滤条件引用的是本节点 mapping 产出的新字段时,先做 mapping、后判条件(否则字段还不存在)
引擎需在构图期分析条件所引用的字段来源,自动决定顺序,不由用户配置。
支持的 mapping 类型:
| 类型 | 语义 |
|---|---|
direct | 从上游直接取字段值 |
constant | 注入常量(long / double / string / bool) |
location | 由 IP 求地理位置;查询失败返回 "unknown" 而非 null |
concat | 多字段拼接 |
与 1.x 的差异:1.x 声明了
concat、bool、variable三种 mapping,但离线引擎未注册,调用返回 null。2.0 要求全部声明的类型均有实现(见 ADR-0005)。
1.3 aggregate —— 有状态聚合
核心节点。在窗口内按 key 维护聚合状态。
- 输入:上游上下文中 key 字段与被聚合字段的值
- 输出:该 key 在当前窗口内的累计聚合值(注意:是累计值,不是本次增量)
- 状态:每个 key 一份聚合状态
groupbykeys 的个数决定状态的形态:
| 个数 | 形态 | 示例 |
|---|---|---|
| 0 | 全局单值 | 全站本小时订单总数 |
| 1 | 一级 key | 每个 IP 的请求数 |
| 2 | 二级 key | 每个 IP 下每个页面的请求数 |
1.4 sequence —— 相邻事件求差
- 输入:某个数值字段(实际使用中恒为
timestamp)+ key - 输出:
当前值 − 上一次值 - 状态:每个 key 保存上一次的值
- 首条事件返回"不通过"(无前值可比),不返回 0
典型用途是计算点击间隔,用于识别机器行为(间隔过于规整或过短)。
与 1.x 的差异:1.x 的实现有两个确定性缺陷——多分组键时 key 数组构造循环写成了
get(0)而非get(i),且构造出的 key 数组从未被赋值,导致实际以全 null 数组查询状态。2.0 修正,多键行为符合定义。此外 1.x 只实现了减法,2.0 保持一致(其他运算无实际用例)。
1.5 dual —— 双变量二元运算
- 输入:两个上游变量各自的
value - 输出:
第一个 op 第二个 - key 取自第二个变量
- 任一侧为 null 时不通过
支持的运算:+、-、*、/。除法的除数为 0 时返回 null(不通过),不返回 Infinity 或抛异常。
跨窗口合并时必须先合并两侧的原始值,再做运算,而不是对已算出的结果做运算。例如合并两个小时的"登录失败率",正确做法是分别合并失败数与总数,再相除。
dual 的值是计算时刻的快照
这是一处容易被忽略、但必须写明的时间语义差异:
| 查询时的行为 | |
|---|---|
aggregate | 按查询时刻重算 —— 滑动窗口会先裁掉已过期的事件 |
dual / sequence / top | 返回上次计算时存下的值 —— 不随查询时刻变化 |
后果是:同一时刻查询一个 dual 变量和它的上游 aggregate,两者可能不一致。例如上游"5 分钟内动态请求数"在最后一次事件到达时是 7,之后有事件过期、查询时重算为 6,而 dual 变量仍返回基于 7 算出的结果。
这是有意的,而非缺陷:dual 表达的是"事件发生那一刻两个指标的关系",把它设计成查询时重算会带来两个问题——需要为每次查询重建上游的完整状态,代价高;而且策略判定关心的本来就是事件发生瞬间的比值,不是查询瞬间的。
实现方需要注意的是保持一致:参考实现与生产引擎在这一点上必须相同,否则跨语言对照会在 dual 变量上出现无法解释的差异。这条语义由 tests/golden/vectors/graph-expected.json 的快照守住。
1.6 top —— TopN
- 输入:一个分组计数型的上游变量
- 输出:按值降序排列的前 N 项,每项为
{key, value} - 默认 N = 100,可由
param指定
值相等时的排序:按 key 的字典序升序作为次级排序,保证结果稳定可比。
与 1.x 的差异:1.x 的 TopN 写路径基本未实现(
compute()中对应分支是空块,getCacheWrappers()恒返回空列表),实际依赖上游分组计数节点的数据现算现排,且未定义值相等时的顺序,结果不稳定。2.0 明确定义。
二、聚合算子
下表中"空窗口"指该 key 在窗口内没有任何有效输入时的返回值。
2.1 计数类
| 算子 | 输入类型 | 输出类型 | 语义 | 空窗口 |
|---|---|---|---|---|
count | 任意 | long | 有效输入的条数(null 不计) | 0 |
distinct_count | long / double / string | long | 去重后的基数,见 §2.5 | 0 |
group_count | 任意 + param 指定的分组字段 | map⟨string, long⟩ | 按分组字段再分组的计数 | 空 map |
2.2 数值类
| 算子 | 输入类型 | 输出类型 | 语义 | 空窗口 |
|---|---|---|---|---|
sum | long / double | 同输入 | 求和 | 0 |
max / min | long / double | 同输入 | 最大 / 最小值 | null |
avg | long / double | double | 算术平均 | null |
variance | long / double | double | 样本方差,分母为 n−1 | n ≤ 1 时返回 0.0 |
stddev | long / double | double | 样本标准差 = sqrt(variance) | n ≤ 1 时返回 0.0 |
cv | long / double | double | 变异系数 = stddev / avg | avg 为 0 或 n ≤ 1 时返回 null |
group_sum | long / double + 分组字段 | map⟨string, 数值⟩ | 按分组字段求和 | 空 map |
⚠️ 与 1.x 的差异:
stddev曾是方差1.x 中名为
stddev的算子实际返回方差,没有开平方,计算式为 ;而变异系数cv内部才做了sqrt。2.0 更正:
stddev返回真正的标准差,并新增独立的variance算子。迁移影响:使用
stddev的变量,新值约为旧值的平方根,量级差异明显。迁移工具默认把 1.x 的stddev变量映射为 2.0 的variance以保持结果一致;若要改用标准差,需手动调整并重新校准依赖它的策略阈值。
2.3 取值类
| 算子 | 输入类型 | 输出类型 | 语义 | 空窗口 |
|---|---|---|---|---|
first | 任意 | 同输入 | 窗口内事件时间最早的一条的值 | null |
last | 任意 | 同输入 | 窗口内事件时间最晚的一条的值 | null |
lastn | 任意 + param=N | list⟨同输入⟩ | 最近 N 条的值,按时间倒序(最新在前) | 空列表 |
distinct | 任意 | list⟨同输入⟩ | 去重后的值集合,按首次出现顺序 | 空列表 |
collection | 任意 | list⟨同输入⟩ | 全部值,保持到达顺序 | 空列表 |
last_value | map | 同 map 的 value 类型 | 取 map 型上游变量的最新值 | null |
global_latest | 任意 | 同输入 | 全局最新值(不分 key) | null |
first / last 依据事件时间而非到达顺序。 这是与 1.x 的隐含差异:1.x 无事件时间语义,实际是到达顺序;2.0 在乱序到达时结果不同(且更正确)。
lastn 是长期画像的主力算子,例如"账号最近 10 个登录 IP"。当 N 超过窗口内实际条数时返回全部,不补空位。
时间戳相同时按到达顺序倒序。 这条容易被忽略:若只按时间排序而不定义同值时的次序,不同实现会给出相反的结果 —— 稳定排序会保留插入顺序,从而返回最早的 N 条,与"最近 N 条"的语义正好相反。
这一条是在跨语言对照中发现的:两套实现对同时间戳的处理相反,而当时的共享向量里没有一条用例的时间戳相同,因此从未暴露。已补入向量
lastn-equal-timestamps。
2.4 合并类
| 算子 | 输入类型 | 输出类型 | 语义 |
|---|---|---|---|
merge | map | map | 合并多个 map,键冲突时取较新的值 |
merge_value | map | map | 合并多个 map,键冲突时对值做求和 |
top / topn | map⟨string, 数值⟩ | list⟨{key, value}⟩ | 见 §1.6 |
2.5 去重计数的精度模式
这是 2.0 与 1.x 差异最大、也最需要注意的算子。
2.0 的行为
由变量的 function.config.distinct_mode 控制:
| 模式 | 行为 |
|---|---|
exact(默认) | 精确去重。基数超过 approx_threshold(默认 100000)时自动降级为近似,并在结果元数据中标记 approximate: true |
approx | 始终使用 HyperLogLog,参数 log2m = 14(16384 个 register),标准误差约 0.8% |
选择精确作为默认值,是因为风控场景中绝大多数去重计数的基数不大(单个 IP 一小时内关联的设备数、单个账号的登录城市数),精确计算的成本可以接受,而精度对策略阈值的稳定性更重要。
1.x 的行为(不兼容)
1.x 采用"前 20 个精确 + 溢出走 HLL"的混合结构:
- 前 20 个不同值用定长哈希集合精确保存,哈希函数为
String.hashCode(),线性探测最多 3 次,3 次冲突即放弃该值 - 超过 20 个后转入 HyperLogLog,参数
log2m = 9、regwidth = 5(512 个 register),标准误差约 4.6%,此时哈希函数改用murmur3_32 - 最终结果为
HLL 基数 + 20
此外,1.x 在跨小时合并时对去重集合只做列表拼接、不做去重,导致系统性高估。
迁移影响
全部 distinct_count 类变量(1.x 内置资产中约 60 个)的历史值与 2.0 的新值不可直接比较,通常 2.0 的值会略低(消除了高估与 20 的固定偏移)。依赖这些变量的策略阈值需要重新校准,建议并行运行 1~2 周后再调整。
三、过滤条件算子
2.0 要求下表全部算子均有引擎实现。 1.x 中配置层声明了完整算子集,但离线引擎仅实现了字符串的 contains / == / !=,数值的 > / >= / < / ==(Double 仅 <),IP 地理类算子的处理器直接返回 null 完全不生效——用户配得出来的变量引擎跑不了。这是 2.0 引入 schema 强制校验的直接动因。
3.1 通用
| 算子 | 适用类型 | 语义 |
|---|---|---|
== / != | 全部 | 相等 / 不等 |
empty / !empty | string | 为空串或 null / 反之 |
3.2 数值与时间
| 算子 | 适用类型 | 语义 |
|---|---|---|
> >= < <= | long / double | 数值比较 |
类型严格:输入值的实际类型与声明类型不符时,条件判定为不通过,而不是尝试转换。
3.3 字符串
| 算子 | 语义 |
|---|---|
contains / !contains | 包含子串 |
startwith / !startwith | 前缀匹配 |
endwith / !endwith | 后缀匹配 |
regex / !regex | 正则匹配(整串匹配语义) |
in / !in | 属于给定集合,集合以逗号分隔 |
containsby / !containsby | 反向包含:字段值是给定值的子串 |
正则表达式在配置保存时编译校验,非法正则拒绝保存。引擎侧对单次匹配设置执行上限,防止灾难性回溯。
3.4 IP 地理位置
| 算子 | 语义 |
|---|---|
locationequals / !locationequals | IP 归属地等于给定地名 |
locationcontainsby / !locationcontainsby | IP 归属地包含于给定地名集合 |
地理库查询失败时,归属地按 "unknown" 参与比较,不使条件报错。
3.5 复合条件
and / or / not 可任意嵌套。求值采用短路语义:and 遇到第一个 false 即返回,or 遇到第一个 true 即返回。not 只取第一个子条件。
四、时间窗口
4.1 窗口类型
| period.type | 窗口语义 |
|---|---|
last_n_seconds | 滑动窗口,长度为 value 秒 |
hourly | 滚动窗口,整点对齐,长度为 value 小时 |
last_n_hours / last_n_days | 滑动窗口 |
today | 自然日窗口,按部署时区的零点对齐 |
ever | 无界,由数据保留期约束 |
self | 无窗口,对当前值直接变形 |
4.2 事件时间与迟到数据
窗口按事件时间划分,不是处理时间。
迟到事件的处理:在 allowedLateness(默认 60 秒,可配)之内到达的迟到事件仍会更新对应窗口的结果;超出该范围的事件进入侧输出流,记录指标并可选落盘供审计,不静默丢弃。
水位线是流级属性,不按 key 维护。 一条事件是否迟到,取决于它的事件时间与整个输入流当前水位线的关系,而不是与它所属 key 的历史最大时间戳的关系。
这一条容易被实现者忽略,后果是:如果按 key 维护水位线,那么每个新出现的 key 的首个事件都不可能被判定为迟到——哪怕它的事件时间比流里其他数据晚了几小时。攻击者可以借此用不断变化的 key(新 IP、新设备号)绕过迟到检测。
本条规定是在编写参考实现时补入的:实现之初把水位线放在了每个 key 的窗口状态里,迟到检测的测试因此失败,才发现原规格没有说明水位线的作用域。
与 1.x 的差异:这是取消离线重算的前提
1.x 的窗口由数据推进而非定时器:事件时间小于当前窗口的事件直接丢弃且无任何记录;引擎过载时还会主动丢事件;进程重启丢失当前窗口全部状态。正是这三点使得 1.x 必须每小时重放一遍事件日志来重算。
2.0 依托 Flink 的事件时间语义、allowedLateness、反压与 Checkpoint,四个根因均不存在,离线重算层整体取消。详见 ADR-0002。
连带的变化:1.x 每小时 0~2 分存在数据空窗(上一小时在线状态已清、离线结果未产出),2.0 不存在;1.x 查询需按"当前小时查在线、历史小时查离线"分流,2.0 查询路径统一。
4.3 窗口结果的可见性
聚合结果在窗口内持续可见(每条事件更新后即可查询到最新累计值),而非等窗口关闭才输出。这是风控场景的硬要求——策略需要在攻击进行中就能判定,不能等到整点。
五、类型推导
算子的合法性由三元组决定:窗口类型 × 操作数类型 × 算子。三者组合不合法时,变量在保存时即被拒绝,不允许留到运行时报错。
推导规则的完整定义见 类型推导规则。核心约束:
- 数值类算子(
sum、avg、variance、stddev、cv)只接受 long / double max/min接受 long / double,输出类型与输入一致distinct_count接受 long / double / string,输出恒为 long- 无窗口(
self)时不允许使用需要累计状态的算子 - 二元运算的输出类型:
long op long → long,但long / long → double
1.x 的类型推导实现(
variable_function.py中的五张 calculator map 加value_type.py)是整套体系中最完整、最值得原样移植的部分——它保证了"窗口类型 × 操作数类型 × 算子 → 结果类型"的封闭可验证性。2.0 保留其结构,并补齐 1.x 中声明了但引擎未实现的组合。
六、回归测试要求
每个算子必须提供以下测试,缺一则 CI 失败(见 ADR-0005):
- 正常路径:典型输入的正确结果
- 空窗口:符合本文规定的返回值
- null 输入:确认被跳过而非当作零值
- 类型不符:确认按本文规定拒绝或不通过
- 窗口边界:窗口切换时状态正确重置
- 迟到数据:allowedLateness 之内正确更新,之外进入侧输出
此外,tests/golden/ 下维护与 1.x 的对照用例:同一批合成事件分别喂给 1.x 与 2.0,逐变量比对。差异必须能被本文标注的语义变更解释,不能解释的差异即为缺陷。