算子语义规格

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 —— 事件解包

计算图的根节点,每张图有且只有一个。

  • 输入:一条原始事件
  • 输出:事件上下文,包含事件的全部属性,以及 appnamekeytimestampvalue 五个固定字段
  • 状态:无
  • 恒通过:不做任何过滤

1.2 filter —— 过滤与字段派生

  • 输入:上游上下文
  • 输出:不通过时剪枝;通过时把 mapping 产出的派生字段写入上下文,并透传 valuetimestamp
  • 状态:无

求值顺序是一个必须精确保留的语义:

  • 当过滤条件引用的是上游已有字段时,先判条件、后做 mapping(条件不满足则跳过 mapping 的计算开销)
  • 当过滤条件引用的是本节点 mapping 产出的新字段时,先做 mapping、后判条件(否则字段还不存在)

引擎需在构图期分析条件所引用的字段来源,自动决定顺序,不由用户配置。

支持的 mapping 类型:

类型语义
direct从上游直接取字段值
constant注入常量(long / double / string / bool)
location由 IP 求地理位置;查询失败返回 "unknown" 而非 null
concat多字段拼接

与 1.x 的差异:1.x 声明了 concatboolvariable 三种 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_countlong / double / stringlong去重后的基数,见 §2.50
group_count任意 + param 指定的分组字段map⟨string, long⟩按分组字段再分组的计数空 map

2.2 数值类

算子输入类型输出类型语义空窗口
sumlong / double同输入求和0
max / minlong / double同输入最大 / 最小值null
avglong / doubledouble算术平均null
variancelong / doubledouble样本方差,分母为 n−1n ≤ 1 时返回 0.0
stddevlong / doubledouble样本标准差 = sqrt(variance)n ≤ 1 时返回 0.0
cvlong / doubledouble变异系数 = stddev / avgavg 为 0 或 n ≤ 1 时返回 null
group_sumlong / double + 分组字段map⟨string, 数值⟩按分组字段求和空 map

⚠️ 与 1.x 的差异:stddev 曾是方差

1.x 中名为 stddev 的算子实际返回方差,没有开平方,计算式为 (squareSumsum×avg)/(count1)(\text{squareSum} − \text{sum} \times \text{avg}) / (\text{count} − 1);而变异系数 cv 内部才做了 sqrt

2.0 更正:stddev 返回真正的标准差,并新增独立的 variance 算子。

迁移影响:使用 stddev 的变量,新值约为旧值的平方根,量级差异明显。迁移工具默认把 1.x 的 stddev 变量映射为 2.0 的 variance 以保持结果一致;若要改用标准差,需手动调整并重新校准依赖它的策略阈值。

2.3 取值类

算子输入类型输出类型语义空窗口
first任意同输入窗口内事件时间最早的一条的值null
last任意同输入窗口内事件时间最晚的一条的值null
lastn任意 + param=Nlist⟨同输入⟩最近 N 条的值,按时间倒序(最新在前)空列表
distinct任意list⟨同输入⟩去重后的值集合,按首次出现顺序空列表
collection任意list⟨同输入⟩全部值,保持到达顺序空列表
last_valuemap同 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 合并类

算子输入类型输出类型语义
mergemapmap合并多个 map,键冲突时取较新的值
merge_valuemapmap合并多个 map,键冲突时对值做求和
top / topnmap⟨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 = 9regwidth = 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 / !emptystring为空串或 null / 反之

3.2 数值与时间

算子适用类型语义
> >= < <=long / double数值比较

类型严格:输入值的实际类型与声明类型不符时,条件判定为不通过,而不是尝试转换。

3.3 字符串

算子语义
contains / !contains包含子串
startwith / !startwith前缀匹配
endwith / !endwith后缀匹配
regex / !regex正则匹配(整串匹配语义)
in / !in属于给定集合,集合以逗号分隔
containsby / !containsby反向包含:字段值是给定值的子串

正则表达式在配置保存时编译校验,非法正则拒绝保存。引擎侧对单次匹配设置执行上限,防止灾难性回溯。

3.4 IP 地理位置

算子语义
locationequals / !locationequalsIP 归属地等于给定地名
locationcontainsby / !locationcontainsbyIP 归属地包含于给定地名集合

地理库查询失败时,归属地按 "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 窗口结果的可见性

聚合结果在窗口内持续可见(每条事件更新后即可查询到最新累计值),而非等窗口关闭才输出。这是风控场景的硬要求——策略需要在攻击进行中就能判定,不能等到整点。


五、类型推导

算子的合法性由三元组决定:窗口类型 × 操作数类型 × 算子。三者组合不合法时,变量在保存时即被拒绝,不允许留到运行时报错。

推导规则的完整定义见 类型推导规则。核心约束:

  • 数值类算子(sumavgvariancestddevcv)只接受 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):

  1. 正常路径:典型输入的正确结果
  2. 空窗口:符合本文规定的返回值
  3. null 输入:确认被跳过而非当作零值
  4. 类型不符:确认按本文规定拒绝或不通过
  5. 窗口边界:窗口切换时状态正确重置
  6. 迟到数据:allowedLateness 之内正确更新,之外进入侧输出

此外,tests/golden/ 下维护与 1.x 的对照用例:同一批合成事件分别喂给 1.x 与 2.0,逐变量比对。差异必须能被本文标注的语义变更解释,不能解释的差异即为缺陷。