多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

RisingWave 流式聚合内部实现解析:HashAggExecutor 的 AggCalls、AggState、AggGroups 与持久化机制

RisingWave 流式聚合内部实现解析:HashAggExecutor 的 AggCalls、AggState、AggGroups 与持久化机制 数据库流处理后端数据工程【免费下载链接】risingwaveEvent streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.项目地址https://gitcode.com/gh_mirrors/ri/risingwave点击查看免费下载本文是 RisingWave 开发者文档 aggregation.md 的深度展开版聚焦流式执行引擎中 HashAggExecutor 的内部实现。你将了解到 HashAggExecutor 的四大核心组件AggCalls、AggState、AggGroups、Persisted State如何协同完成分组聚合值型聚合与 MaterializedInput 型聚合在状态存储上的本质区别以及 AggGroups 的初始化、缓存与持久化流程。读完本文你既能看懂聚合执行器的整体架构也能据此在 hash_agg.rs、agg_group.rs 等源码文件中继续深入。一、聚合在 RisingWave 流式执行中的定位在 RisingWave 中流式聚合是流处理的核心算子之一。与批处理聚合一次性扫描全部输入不同流式聚合需要增量地消费持续到达的数据流在内存中维护每个分组的中间状态并在 checkpoint屏障 barrier到来时把状态持久化到下层状态存储中以保证故障恢复后结果不丢失、不重复。从执行器的分类看聚合类执行器主要有HashAggExecutor处理带GROUP BY的分组聚合也是本文的主角SimpleAggExecutor处理不带GROUP BY的全局聚合group key 为None其实现位于 simple_agg.rsstateless_simple_agg.rs则覆盖max/min等基于 state table 排序键直接求值的无状态场景。所有聚合执行器共用一套AggState/AggGroup/AggStateStorage抽象源码集中在 src/stream/src/executor/aggregate 目录。注意原设计文档中## Frontend与## Expression Framework两节仍标注为 TODO即前端优化器如何生成 AggCall 计划与表达式框架聚合函数如何求值尚未有正式设计文档。本文与当前文档一致只深入已完成的## HashAggExecutor一节。二、HashAggExecutor 概览四大核心组件HashAggExecutor 内部由 AggCalls、AggState、AggGroups 与持久化状态Persisted State四部分组成。在HashAggExecutor内部聚合功能被拆分为 4 个主要组件AggCalls查询中的聚合调用。例如查询包含SUM(v1)、COUNT(v2)时AggCalls 就是SUM和COUNT。AggState用于计算聚合调用输出结果的中间状态。在每个聚合组AggGroup内部每个 AggCall 都对应一份独立的 AggState。AggGroups按聚合组维度创建的执行单元。例如GROUP BY x1, x2时每一个(x1, x2)唯一取值组合都对应一个独立的组。Persisted State持久化状态。当屏障barrier到来时内存中的聚合状态会被写入下层状态表。2.1 数据流与屏障驱动的处理循环从 hash_agg.rs 的源码注释与实现可以看出HashAggExecutor的整体工作方式如下从上游input拉取数据流将到达的数据块StreamChunk应用到对应的聚合状态上处理过程中通过dirty_groups对应源码注释中的group_change_set记录本 epoch 内被修改过的组收到屏障时调用状态后端的flush把所有修改刷入状态存储同时遍历dirty_groups基于状态变化生成并向下游输出结果数据块。这个增量更新 屏障刷盘 变更输出的循环是理解流式聚合的关键聚合结果不是等输入结束后一次性计算而是每个 epoch 输出一次增量变更Insert/Update/Delete。主循环定义在execute_inner中对三类消息分别处理Message::Chunk调用apply_chunk将数据应用到各组的聚合状态若dirty_groups的预估堆大小超过max_dirty_groups_heap_size会提前触发flush_data以防止 OOM内存溢出Message::Watermark缓冲组键列上的水位线用于 EOWC 模式或下游水位传播Message::Barrier执行flush_data输出本 epoch 的结果变更、提交所有状态表commit_state_tables、刷新指标并转发屏障。2.2 构造参数HashAggExecutor 的输入装配执行器由HashAggExecutor::new(args)构造参数封装在AggExecutorArgs中见 mod.rsagg_calls/row_count_index聚合调用列表以及count(*)在列表中的下标用于判断某组是否有输入行storages每个 agg call 对应的状态存储方式Value或MaterializedInputintermediate_state_table值型聚合共享的中间状态表distinct_dedup_tables为去重 DISTINCT 聚合列而维护的独立状态表一个去重列一张可被多个 agg call 共享extreme_cache_size极值类聚合min/max状态缓存上限HashAggExecutorExtraArgs包含group_key_indices分组键在输入中的列索引、chunk_size单次输出数据块最大行数、max_dirty_groups_heap_size脏组堆大小上限、emit_on_window_close是否采用窗口关闭时发射策略。一个值得注意的构造细节源码假定状态表主键的前缀恰好就是分组键group_key_table_pk_projection并在构造时用断言校验——这是 AggGroup 能从状态表按分组键前缀高效恢复状态的基石。三、AggCalls聚合调用的描述与执行AggCalls 是查询中声明的聚合调用集合比如SUM(v1)、COUNT(v2)各对应一个 AggCall。它是要算什么的声明层真正怎么算由底层聚合函数实现。在HashAggExecutor内部agg_calls: VecAggCall保存聚合调用的静态描述聚合类型、入参列、过滤条件、是否 DISTINCT 等agg_funcs: VecBoxedAggregateFunction是通过build_retractable构建的可撤销retractable聚合函数实例——每个 AggCall 一个。所谓可撤销是指这类聚合函数既能增量吸收新行update也能在收到删除DELETE / - 操作时回退状态这是流式聚合正确处理变更流的前提。在apply_chunk阶段每个 AggCall 还会计算自己的行可见性call visibility对min/max/string_agg这类函数会额外屏蔽 NULL 值行agg_call_filter_res见 mod.rs对带FILTER子句的聚合会求值过滤表达式得到新的可见性位图。随后数据按可见性投影分别喂给对应的 AggState。四、AggState单值状态与物化输入状态AggState 是单个聚合调用用于计算输出的状态。在每一个聚合组内部每个 AggCall 都有自己的一份 AggStatestates: VecAggState与agg_calls一一对应。源码 agg_state.rs 中AggState是如下枚举pub enum AggState { /// 单个标量值状态例如 count、sum、append-only 的 min/max Value(AggregateState), /// 物化输入块状态例如非 append-only 的 min/max、string_agg MaterializedInput(BoxMaterializedInputState), }Value(AggregateState)聚合状态就是一个标量值或极小的结构体。count、sum以及仅追加场景下的min/max都属于此类——因为输入只增不减一个值足以表达全部历史。创建时若能从编码状态encoded_state解码则恢复否则用func.create_state()新建。MaterializedInput(BoxMaterializedInputState)聚合状态是输入行的物化集合。非 append-only 的min/max输入可能删除旧值、string_agg需要按序拼接字符串属于此类因为单靠一个标量无法在删除发生时还原正确的输出。4.1 AggState 的两个关键操作AggState提供两个核心方法apply_chunk(chunk, call, func, visibility)把输入块应用到状态上。Value状态会按call.args.val_indices()投影出参数列后调用func.update(state, chunk)MaterializedInput状态则把整个块不投影交给MaterializedInputState::apply_chunk管理。get_output(storage, func, group_key)产出该聚合调用的当前输出。Value状态直接调用func.get_result(state)MaterializedInput状态则需要结合状态表读取来得出结果见下文输出构建一节。此外还有reset用于把值型状态重置回初始状态——这在分组行数降为 0、需要输出 NULL/0 语义时至关重要。五、AggStateStorage值型与 MaterializedInput 型的两套持久化与AggState内存态对应AggStateStorage描述聚合状态的持久化方式pub enum AggStateStorageS: StateStore { /// 状态以单个值的形式存放在中间状态表中 Value, /// 状态以输入块物化的形式存放在独立状态表中 MaterializedInput { table: StateTableS, // 独立的状态表 mapping: StateTableColumnMapping, // 状态表列与输入块的列映射 order_columns: VecColumnOrder, // 排序键的列索引与排序方向 }, }5.1 值型聚合共享一张中间状态表对于 value 类型聚合整个执行器只有一张intermediate_state_table所有组的所有值型聚合状态以一行一组的方式存储——这一行由分组键前缀 每个值型聚合的编码状态列构成。GroupKey见 agg_group.rs封装了这一映射关系row_prefix具体分组键值table_pk_projection分组键到状态表主键的投影table_row()取状态表行的前缀用于写行table_pk()取状态表主键前缀用于按组查询。该设计假定状态表主键前缀即分组键HashAggExecutor::new中有断言校验因此按组恢复状态时只需一次主键前缀范围查询。5.2 MaterializedInput 型聚合每个聚合一张独立状态表对于 MaterializedInput 类型聚合每个聚合调用各有一张状态表AggStateStorage::MaterializedInput表内按分组存放该聚合需要的输入状态。非 append-only 的 min/max、string_agg 都需要把输入行或其关键列物化保存才能在收到删除时正确回退。MaterializedInputStateminput.rs内部维护一个类型化状态缓存Boxdyn AggStateCache对min/max这类极值聚合使用TopNStateCache容量受extreme_cache_size限制对string_agg等使用有序状态缓存OrderedStateCache。写入状态表的时机同样由屏障驱动输入块会先写入tablewrite_chunk并在flush时真正落盘。5.3 屏障时的持久化与提交在flush_data阶段每个脏组依次执行两个动作build_states_change对比上一轮中间状态与当前中间状态计算出差量记录Insert/Update/Delete写入intermediate_state_tableEOWC 模式则先放入SortBuffer缓冲MaterializedInput 状态自身不需要写入这张表它们的编码列恒为None因为其状态已经物化在各自的独立状态表中build_outputs_change基于当前状态计算聚合输出与上一轮输出对比生成变更行交给下游。最后在屏障处理中commit_state_tables会对所有状态表各 MaterializedInput 表、distinct 去重表、中间状态表执行 epoch 提交。六、AggGroups按组管理聚合状态AggGroups 按聚合组维度创建。GROUP BY x1, x2时每一个(x1, x2)唯一组合对应一个组AggGroup负责该组内所有 AggCall 的状态管理与变更输出。AggGroupS, Strtgagg_group.rs的关键字段ctx组上下文含分组键group_key: OptionGroupKey。SimpleAggExecutor中为Nonestates: VecAggState本组内每个 AggCall 对应的聚合状态prev_inter_states上一轮写入中间状态表的编码中间状态prev_outputs上一轮输出给下游的结果行row_count_indexcount(*)在调用列表中的下标emit_on_window_close是否为窗口关闭发射模式。6.1 两种输出策略AlwaysOutput 与 OnlyOutputIfHasInputAggGroup通过泛型策略参数Strtg: Strategy决定何时向下游输出变更AlwaysOutput无论该组是否有输入行都输出聚合结果用于支持approx_percentile这类需要每个 epoch 都输出更新、以配合row_merge无状态执行器的场景。首次构建时输出Insert之后恒为Update即使更新是无变化的no-op也照样输出OnlyOutputIfHasInput仅在组内有输入行时输出。用count(*)判断组内行数行数0→0不输出0→n输出Insertn→0输出Deleten→m且结果有变化才输出Update。HashAggExecutor 使用OnlyOutputIfHasInputtype AggGroupS GenericAggGroupS, OnlyOutputIfHasInput。6.2 输出构建的细节get_outputs是输出构建的核心若当前组行数为 0会先reset所有值型状态——因为某些聚合如sum在没有输入行时应输出 NULL另一些如sum0应输出 0重置才能保证语义正确。然后并行地对每个状态调用state.get_output(...)MaterializedInput 状态可能涉及状态表读取所以build_outputs_change比build_states_change更耗时。对行数为负的异常场景上游传来不一致的 DELETE源码会做防御性处理将行数视为 0 并抑制下游输出避免因不存在的键导致下游 panic对应 issue 14031。七、AggGroups 的初始化与 AggGroupCacheAggGroup 初始化发生在 AggGroupCache 未命中时要么缓存被逐出要么这是一个全新的分组键。AggGroups 的初始化发生在对应的聚合组在AggGroupCache中找不到时。触发原因有两类AggGroupCache中的条目被逐出LRU 缓存容量受限这是一个全新的分组键此前从未出现过。7.1 AggGroupCache 的作用初始化一个 AggGroup 是相对昂贵的操作——需要从中间状态表读取上一轮编码状态、为每个 AggCall 创建/恢复 AggState、并在非 EOWC 模式下预先计算prev_outputs。因此执行器用ManagedLruCache缓存它们type AggGroupCacheK, S ManagedLruCacheK, OptionBoxedAggGroupS, PrecomputedBuildHasher;ExecutionVars中维护两个关键集合agg_group_cacheLRU 缓存键为哈希键HashKey值为可空的 AggGroupdirty_groups本 epoch 被修改过的组EstimatedHashMap脏组在屏障时被刷写。7.2 touch_agg_groups查找、恢复与新建touch_agg_groupshash_agg.rs是初始化逻辑的入口对每个出现的分组键若该键已在dirty_groups中直接跳过否则尝试agg_group_cache.get_mut(key)缓存命中把 AggGroup 从缓存搬移到dirty_groupsagg_group.take()留下None占位无需重新创建缓存未命中lookup miss并发buffer_unordered(10)执行AggGroup::create(...)——它会先get_row从intermediate_state_table按分组键主键前缀读取上一轮状态若存在则恢复AggState::create用agg_func.decode_state解码若不存在则用create_state新建空状态。创建完成后AggGroup被放入dirty_groups等待本轮应用数据与屏障刷写。由于可能同时初始化大量组源码注释特别说明collect_vec是为了规避缓存借用导致的生存期问题。AggGroup::create还会统计状态缓存命中/缺失AggStateCacheStats供指标上报。7.3 脏组刷写、缓存回填与指标屏障或脏组堆过大时执行flush_data遍历dirty_groupsbuild_states_change写中间状态表、build_outputs_change生成输出块随后把所有脏组放回agg_group_cachedirty_groups.drain()cache.put并evict()驱逐超出容量的条目维持缓存水位。执行器还通过ExecutionStats记录并上报一组缓存指标agg_lookup_miss_count/agg_total_lookup_count组粒度命中率、agg_chunk_lookup_miss_count/agg_chunk_total_lookup_count数据块粒度命中率、agg_cached_entry_count缓存条目数、agg_state_cache_lookup_count/agg_state_cache_miss_countAggState 内部缓存命中率。这些指标可以在 Grafana 中观察聚合执行器的缓存效率。7.4 与 EOWC窗口关闭发射的关系当emit_on_window_close为 true如时间窗口聚合采用窗口关闭时才发射策略时flush_data走另一条路径状态变更先写入SortBuffer按窗口列排序的缓冲见 hash_agg.rs 中的buffer: SortBufferS收到窗口列水位线后buffer.consume(watermark, ...)把水位线以下的组逐行取出用AggGroup::for_eowc_output重建输出组并产出结果EOWC 下每组只输出一次恒为Insertvnode 位图变化时还需refill_cache重建缓冲缓存避免输出不属于本 actor 的行。八、写在最后从文档到源码的阅读路径原设计文档 aggregation.md 以简练的篇幅勾勒了 HashAggExecutor 的四大组件与 AggGroup 初始化逻辑本文在此基础上补充了完整的源码证据。如果你想继续深挖建议按以下顺序阅读hash_agg.rs执行器主循环、apply_chunk/flush_data/touch_agg_groups的完整实现agg_group.rsAggGroup、GroupKey与两种输出策略agg_state.rsAggState与AggStateStorage的两种形态minput.rsMaterializedInput 状态的缓存与状态表交互mod.rs执行器公共构造参数与可见性计算simple_agg.rs无分组键的全局聚合如何复用同一套AggGroup抽象。此外integration_tests.rs 与 agg_executor.rs 提供了聚合执行器的测试装配方式是理解参数如何组装成一个可运行的聚合执行器的绝佳补充。对于分布式一致性还需结合commit_state_tables与屏障提交机制理解 epoch 对聚合状态的保护语义。赞分享数据库流处理后端数据工程【免费下载链接】risingwaveEvent streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.项目地址https://gitcode.com/gh_mirrors/ri/risingwave点击查看免费下载相关推荐SuperPlane Factory内部机制揭秘Work Order状态机与Line派发流程的持久化实现SuperPlane Factory内部机制揭秘Work Order状态机与Line派发流程的持久化实现 SuperPlane 是一个开源的「一键式工程工厂」Spring AI 依赖管理实战3 行配置导入 BOM告别版本冲突Spring AI 依赖管理实战3 行配置导入 BOM告别版本冲突 今天加 spring ai openai 过几天加 spring ai pgvecto人工智能大模型AI AgentRAG后端工具调用MCP 服务MCP ClientsRedis持久化机制深度解析RDB与AOF实现原理Redis持久化机制深度解析RDB与AOF实现原理 本文深入解析Redis的两种持久化机制RDB和AOF。RDB通过生成二进制快照文件实现数据持久化具有紧数据库KV存储缓存上一篇如何用Upscayl将模糊图片变高清免费AI超分辨率实战指南下一篇如何使用Typo.css快速实现专业级中文网页排版新手入门指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表