首页/新闻资讯/正文详情

RisingWave Stream Engine Executor 测试编写指南:基于 expect_test 与集成测试的最佳实践

发布时间:2026/9/25 7:58:50 来源:云帆数科 栏目:资讯中心
RisingWave Stream Engine Executor 测试编写指南:基于 expect_test 与集成测试的最佳实践
数据库流处理后端数据工程【免费下载链接】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 的 Stream Engine流式引擎是其物化视图实时刷新能力的核心由大量流式 executor如 Hash Agg、Hop Window、Over Window 等构成。本文聚焦于 src/stream/README.md 所确立的两条测试规范——使用expect_test编写快照测试、以集成测试取代单元测试——结合仓库源码与真实用例完整讲解如何为新的流式 executor 编写可自动更新期望输出、可验证恢复语义的高质量测试。读完本文你将掌握check_until_pending、check_with_script、MockSource等核心测试工具的用法与底层实现并能在 src/stream/tests/integration_tests 中快速落地自己的测试。一、Stream Engine 背景与本文测试对象的定位RisingWave 通过物化视图MV提供实时分析能力所有物化视图都会根据最近的数据更新自动刷新而这一刷新过程由流式引擎承担。关于引擎的整体架构frontend 服务层、compute node 上的长驻 actor、meta service 协调、基于云对象存储的共享状态存储等docs/dev/src/design/streaming-overview.md 有系统性的论述其中包含了流式引擎三条核心设计原则基于 actor 模型的执行引擎每个 actor 独立响应输入消息数据更新与控制信号构成高并发流式引擎状态共享存储状态存储基于云对象存储当前为 AWS S3以获得计算弹性、廉价且近乎无限的容量以及配置变更时的简单性万物皆表、万物皆状态内部存储中的每个对象既是一个逻辑表也是一个内部状态从而可以被 catalog 统一管理并在统一的一致性保证下由流式引擎更新。在 actor 内部executor 是增量计算delta computation的基本单元每个 executor 对应一个关系算子含基表当任一基表收到更新时流式引擎从叶子到根递归计算每个物化视图的变化每个节点收到子节点的更新后计算本地更新并向上游传播。只要保证每个 executor 的正确性就能组合出任意 SQL 查询的增量维护框架。因此executor 的测试质量直接决定了整个流式引擎的正确性——这正是本文测试指南的出发点。二、测试编写总纲两条核心推荐src/stream/README.md 为 executor 测试确立了两条明确的规范使用expect_test编写新测试expect_testrust-analyzer 官方出品的快照测试库可以自动更新期望输出将测试维护成本降到最低。详见check_until_pending的文档及其用法。优先编写集成测试而非单元测试新测试应写在tests/目录即src/stream/tests/integration_tests/而不是src/目录内。这条建议来自 upstream issue #9878 的讨论结论——集成测试更贴近真实运行路径便于验证 executor 在完整消息流下的行为且不污染生产代码树。expect-test已作为 dev-dependency 声明在 src/stream/Cargo.toml 中expect-test 1因此在该 crate 内编写测试无需额外引入依赖。三、集成测试目录结构与基建src/stream/tests/integration_tests/是 Stream Engine executor 集成测试的统一入口。其结构如下main.rs集成测试的 crate 入口声明了各测试模块并定义了统一的prelude模块将expect_test::expect与risingwave_stream::executor::test_utils::prelude::*重导出同时注入本 crate 的snapshot::*工具使得各测试文件可以极简地使用check_until_pending、SnapshotOptions等snapshot.rs快照测试基础设施check_until_pending、check_with_script、SnapshotOptions、脚本 DSL各功能测试模块hash_agg.rsHash 聚合、hop_window.rs滑动窗口、over_window.rs与eowc_over_window.rsOver Window 及其 EOWC 变体、project_set.rs投影集合、materialized_exprs.rs物化表达式。测试模块之间通过 integration_test.toml 定义的 RiseDev 任务串联其中[tasks.do-apply-stream-integration-test] description Apply stream integration test output snapshots category RiseDev - Test dependencies [install-nextest] script #!/usr/bin/env bash set -e UPDATE_EXPECT1 cargo nextest run -p risingwave_stream --test integration_tests --retries 0 echo $(tput setaf 2)Diff applied!$(tput sgr 0) echo Tip: use the alias $(tput setaf 4)./risedev dasit$(tput sgr0). 即执行./risedev dasitalias 为do-apply-stream-integration-test即可一键应用所有集成测试的期望输出。四、核心工具一check_until_pending与SnapshotOptions4.1 基本用法check_until_pending定义于 src/stream/tests/integration_tests/snapshot.rs其职责是驱动 executor 直到其进入 pending 状态然后断言输出与expect完全一致。其 doc 注释给出了新增测试的建议工作流创建 executor 并向其发送输入消息后只需放置这样一行check_until_pending(mut executor, expect![[]]).await;如果内联结果排版不佳也可以改用expect_file!将快照存到独立文件。然后运行带UPDATE_EXPECT1的测试即可自动填充期望值UPDATE_EXPECT1 cargo nextest run -p risingwave_stream # 或 UPDATE_EXPECT1 risedev test -p risingwave_stream该函数的签名与实现如下pub fn check_until_pending( executor: mut BoxedMessageStream, expect: expect_test::Expect, options: SnapshotOptions, ) { let output run_until_pending(executor, mut Store::default(), options); let output serde_yaml::to_string(output).unwrap(); expect.assert_eq(output); }可以看到它内部依赖run_until_pending使用futures::FutureExt::now_or_never()非阻塞地轮询 executor 的try_next()收集所有立即可用的输出消息Chunk、Barrier、Watermark直到 executor 返回None数据流结束或阻塞pending。这一设计的精妙之处在于测试永远不会被卡死——如果 executor 行为异常now_or_never会立即返回 pending测试随即结束并比对输出便于定位问题。4.2SnapshotOptions快照输出的两种控制开关#[derive(Debug, Clone, Default)] pub struct SnapshotOptions { /// 是否对输出 chunk 排序当输出 chunk 无既定顺序时必须开启 pub sort_chunk: bool, /// 是否在每个输出 chunk 之后附带“应用该变更后的结果”可想象为每次输出后执行 SELECT * FROM mv 的结果 pub include_applied_result: bool, }sort_chunk对输出 chunk 调用chunk.sort_rows()后再序列化保证在多行无序输出时快照的确定性见下方 Hash Agg 示例include_applied_result在输出 chunk 后追加应用变更后的整表结果。其底层由Store实现——一个基于BTreeMapDefaultOrderedOwnedRow, usize的多重集apply_chunk依据Op::Insert / UpdateInsert / Delete / UpdateDelete增删行计数最终以DataChunk形式返回“执行SELECT * FROM mv后的完整视图”从而直观验证流式结果的正确性。两个开关均提供 builder 风格方法sort_chunk(bool)、include_applied_result(bool)。4.3 真实用例Hash Agg以 src/stream/tests/integration_tests/hash_agg.rs 中的test_hash_agg_count_sum为例这是最典型的check_until_pending用法创建内存状态存储MemoryStateStore::new()与三列Int64的 schema定义局部 Hash 聚合key 为第 0 列聚合调用为count、sum($1)、sum($2)通过MockSource::channel()取得(MessageSender, MockSource)将 source 转成 executor 后构造 hash agg executor 并execute()依次推送 barrier(1)、含 3 行的 chunk、barrier(2)、含增删混合的 chunk、barrier(3)最后用check_until_pending断言输出快照check_until_pending( mut hash_agg, expect![[r# - !barrier 1 - !chunk |- --------------- | | 1 | 1 | 1 | 1 | | | 2 | 2 | 4 | 4 | --------------- - !barrier 2 - !chunk |- ---------------- | | 3 | 1 | 3 | 3 | | - | 1 | 1 | 1 | 1 | | U- | 2 | 2 | 4 | 4 | | U | 2 | 1 | 2 | 2 | ---------------- - !barrier 3 #]], SnapshotOptions::default().sort_chunk(true), );注意这里开启了sort_chunk(true)因为 Hash 聚合的输出行没有既定顺序排序保证了快照的确定性。快照中U-/U体现了流式引擎的更新语义第二次聚合count 变 1、sum 变 2以“先删旧值、后插新值”的 Update 对呈现这正是增量维护change propagation在输出层的真实形态。同样src/stream/tests/integration_tests/hop_window.rs 中的test_watermark_output_indices2展示了check_until_pending对 watermark 输出的断言- !watermark携带col_idx与val用于验证 hop window 在指定输出列索引下的 watermark 透传与衍生行为。五、核心工具二check_with_script与测试脚本 DSL对于需要多个输入事件、甚至需要模拟恢复recovery的场景check_with_script提供了一种更声明式的方式以 YAML 脚本描述输入事件序列驱动 executor 逐步执行并输出每一步的输入/输出对照快照。其签名与实现要点同样位于 snapshot.rspub async fn check_with_scriptF, Fut( create_executor: F, test_script: str, expect: expect_test::Expect, options: SnapshotOptions, ) where F: Fn() - Fut, Fut: FutureOutput (MessageSender, BoxedMessageStream), { let output executor_snapshot(create_executor, test_script, options).await; expect.assert_eq(output); }create_executor是一个工厂闭包每次调用返回(MessageSender, BoxedMessageStream)——它既用于首次创建 executor也用于脚本中的recovery事件重建 executor 以模拟故障恢复。5.1 脚本 DSL 的输入事件类型脚本由若干SnapshotEvent组成serde_yaml反序列化serde(rename_all lowercase)事件YAML 写法语义Barrier(epoch)- !barrier 1推送指定 epoch 的 barrier内部通过test_epoch构造输出时做对应的右移还原Noop—占位脚本中unreachable!()预留类型Recovery- recovery重新调用create_executor重建(tx, executor)模拟 actor 故障恢复Chunk(String)- !chunk \|2 pretty 格式文本以StreamChunk::from_pretty解析 chunk 文本与快照输出使用同一 pretty 格式便于复制粘贴Watermark { col_idx, val }- !watermark推送指定列索引与值的 watermark当前实现固定为Int64类型5.2 脚本的输出快照结构executor_snapshot会对每个输入事件记录其对应输出事件列表形成Snapshot { input, output: VecSnapshotEvent }结构整体序列化为 YAML。这意味着expect快照中既有输入也有输出任何一步的行为变化都会产生清晰的 diff方便定位是哪个输入引发了回归。此外输出中的 barrier epoch 会执行一次右移barrier.epoch.curr / test_epoch(1)以匹配脚本中人为选定的“小 epoch 数字”保证输入输出一致可比。5.3 真实用例Over Window 与恢复语义src/stream/tests/integration_tests/eowc_over_window.rs 中的test_over_window是脚本化测试的典型范例构造了lag(x, 1)与lead(x, 1)两个窗口函数调用后直接使用脚本驱动check_with_script( || create_executor(calls.clone(), store.clone()), r### - !barrier 1 - !chunk |2 I T I i 1 p1 100 10 1 p1 101 16 4 p2 200 20 - !chunk |2 I T I i 5 p1 102 18 7 p2 201 22 8 p3 300 33 # NOTE: no watermark message here, since watermark(1) was already received - !barrier 2 - recovery - !barrier 2 - !chunk |2 I T I i 10 p1 103 13 12 p2 202 28 13 p3 301 39 - !barrier 3 ###, expect![[r# - input: !barrier 1 output: - !barrier 1 - input: !chunk |- ... output: - !chunk |- ... - input: recovery output: [] ... #]], SnapshotOptions::default(), ) .await;该用例清晰地展示了脚本 DSL 的三大价值chunk 文本直接复用 pretty 格式 1 p1 100 10与输出快照同构读写直观- recovery事件重建 executor且后续输入从“恢复后的状态”继续推进输出中能看到 recovery 后无多余输出output: []且后续窗口计算lag/lead 结果依然正确从而验证了流式引擎的容错恢复语义注释行# NOTE: ...可内嵌在脚本中解释行为意图YAML 解析自动忽略为测试自文档化提供了空间。同样的脚本化模式也应用于 materialized_exprs.rs、over_window.rs 与 hash_agg.rs 中的test_hash_agg_recovery- recovery后复用同一状态存储重建 executor可相互参照。六、配套测试基建MockSource与MessageSender无论选择check_until_pending还是check_with_script输入消息都通过MockSource/MessageSender注入。其定义位于 src/stream/src/executor/test_utils/mock_source.rsMockSource::channel()返回(MessageSender, MockSource)——内部是 tokio 的mpsc::UnboundedChannelMessageMessageSender包装发送端MockSource包装接收端并可into_executor(schema, stream_key)转为标准 executorMessageSender提供了一系列语义化的推送方法方法作用push_chunk(StreamChunk)推送数据变更 chunkpush_barrier(epoch, stop)推送指定 epoch 的 barrierstoptrue时附加停止标记barrier.with_stop()send_barrier(Barrier)直接推送构造好的 barrierpush_barrier_with_prev_epoch_for_test(cur, prev, stop)推送带前序 epoch 的 barrier用于需感知prev_epoch的 executorpush_watermark(col_idx, data_type, val)推送 watermark 消息push_int64_watermark(col_idx, val)便捷推送Int64watermark此外test_utils/mod.rs 中还有一个关键 traitStreamExecutorTestExt其next_unwrap_ready方法可以在不await的情况下取出 executor 的下一条消息若 executor 未就绪则立即 panic避免测试卡死——它是check_until_pending中now_or_never思路的 trait 化补充二者共同保证了“永不阻塞”的测试体验。这些基建都被统一重导出到risingwave_stream::executor::test_utils::prelude集成测试中通过crate::prelude::*见 main.rs即可全部引入测试代码因此非常精简。七、测试的编写与维护工作流总结综合 src/stream/README.md、snapshot.rs 的 doc 与 integration_test.toml为新的 executor 添加测试的完整工作流如下选址将新测试作为集成测试放入src/stream/tests/integration_tests/下的新模块如my_executor.rs并在 main.rs 中注册mod my_executor;不要在src/内写单元测试搭建用MockSource::channel()创建输入通道按目标 executor 的构造参数schema、状态存储、算子配置组装 executor 并调用execute()得到BoxedMessageStream注入输入直接调用tx.push_barrier(...)/tx.push_chunk(StreamChunk::from_pretty(...))或改用check_with_script以 YAML 脚本声明式描述输入含recovery事件放置断言在末尾添加一行check_until_pending(mut executor, expect![[]], SnapshotOptions::default()).await;需要排序时传入.sort_chunk(true)需要整表结果时传入.include_applied_result(true)如内联快照过长可改用expect_file!生成期望值运行UPDATE_EXPECT1 cargo nextest run -p risingwave_stream或仓库内一键命令./risedev dasit自动填充expect!中的快照人工审查仔细检查自动生成的快照是否符合预期语义chunk 的增删改、barrier 透传、watermark 列与值确认无误后提交。这套以快照对比为核心、以自动更新为手段、以集成测试为载体的测试方法论是 RisingWave Stream Engine 长期以来保持 executor 正确性与重构安全性的关键工程实践当你在仓库中修改任一 executor 实现时上述集成测试套件会立即通过快照 diff 揭示任何行为变化从而将回归风险降到最低。八、延伸阅读流式引擎整体架构与设计原则docs/dev/src/design/streaming-overview.md快照测试基础设施check_until_pending/check_with_script/SnapshotOptions/ 脚本 DSLsrc/stream/tests/integration_tests/snapshot.rs集成测试入口与 preludesrc/stream/tests/integration_tests/main.rs测试基建MockSource/MessageSender/StreamExecutorTestExtsrc/stream/src/executor/test_utils/mod.rs、src/stream/src/executor/test_utils/mock_source.rs各功能测试示例Hash Agghash_agg.rs、Hop Windowhop_window.rs、Over Windowover_window.rs 与 eowc_over_window.rs、投影集合project_set.rs、物化表达式materialized_exprs.rs快照自动应用任务与./risedev dasit别名integration_test.toml赞分享数据库流处理后端数据工程【免费下载链接】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点击查看免费下载相关推荐NzbDav与Stremio集成教程通过AIOStreams实现直接Usenet流媒体NzbDav与Stremio集成教程通过AIOStreams实现直接Usenet流媒体 NzbDav是一款功能强大的Usenet流媒体解决方案它结合了WebspaCy 测试体系完全指南基于 pytest 的测试编写、运行与最佳实践spaCy 测试体系完全指南基于 pytest 的测试编写、运行与最佳实践 spaCy 是一套工业级自然语言处理NLP库其核心功能分词、词性标注、依存人工智能NLP机器学习预训练Gutenberg 端到端测试指南基于 Playwright 编写 E2E 测试的最佳实践与源码解析Gutenberg 端到端测试指南基于 Playwright 编写 E2E 测试的最佳实践与源码解析 本篇技术指南以 Gutenberg 项目的官方贡献文档后端前端上一篇Kafka-Map进阶技巧Broker状态监控与性能优化实用指南下一篇Zstandard 压缩格式规范深度解析v0.3.9创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关推荐

wp-calypso 多站点 Dashboard 的包导入限制策略:calypso 与 @automattic 依赖边界实战指南
wp-calypso 多站点 Dashboard 的包导入限制策略:calypso 与 @automattic 依赖边界实战指南

前端CMS 【免费下载链接】wp-calypso The JavaScript and API powered WordPress.com 项目地址: https://gitcode.com/gh_mirrors/wp/wp-calypso 点击查看 免费下载 本文解读 wp-calypso 仓库中新版多站点 Dashboard(服务于 WordPress.com 托管控制台与… · 2026/9/25 7:58:50

CSP-S 2026初赛模拟卷2:考点拆解与备考策略
CSP-S 2026初赛模拟卷2:考点拆解与备考策略

1. 从一份模拟卷说起:CSP-S 初赛到底在考什么如果你正在准备 CSP-S(CCF 非专业级软件能力认证提高组)的第一轮,那你大概率已经刷过不少真题和模拟卷了。但很多人刷题的方式其实很低效——做完对个答案,看看分数&#x… · 2026/9/25 7:58:50

把浏览器交给 AI Agent 之前,BrowserSkill 的 11 项权限与隐私边界怎么验(15 分钟串讲)
把浏览器交给 AI Agent 之前,BrowserSkill 的 11 项权限与隐私边界怎么验(15 分钟串讲)

把浏览器交给 AI Agent 之前,BrowserSkill 的 11 项权限与隐私边界怎么验(15 分钟串讲) 【免费下载链接】BrowserSkill Let AI agents use your real, logged-in browser without interrupting your work. CLI extension for browser automa… · 2026/9/25 7:58:38

PDF图片转Word用什么软件?电脑/网页/手机全覆盖实用攻略
PDF图片转Word用什么软件?电脑/网页/手机全覆盖实用攻略

日常办公、学习中,大家经常遇到一个难题:拿到图片型PDF、扫描件PDF,普通转换根本没用,转完依旧是无法编辑的图片,手动打字费时又费力。这里先科普一个关键知识点:图片类PDF必须依靠OCR文字识别技术&#xf… · 2026/9/25 8:18:56

IT技术岗转网络安全值得吗?成本、路线与就业全景解析
IT技术岗转网络安全值得吗?成本、路线与就业全景解析

我经常在后台收到类似的提问:干了几年IT技术岗,到底要不要转网络安全?说实话,每次看到这种问题,我都能大概猜到提问者的处境——现有工作不算差,但天花板感越来越明显;网络安全听起来热门、有技… · 2026/9/25 8:18:44

Moto 中 Amazon Managed Prometheus(amp)服务的模拟实现与实战指南
Moto 中 Amazon Managed Prometheus(amp)服务的模拟实现与实战指南

Mock测试 【免费下载链接】moto A library that allows you to easily mock out tests based on AWS infrastructure. 项目地址: https://gitcode.com/gh_mirrors/mo/moto 点击查看 免费下载 Amazon Managed Prometheus(AMP,AWS 的托管 Prom… · 2026/9/25 8:18:31

平头哥倚天720/730/750三代CPU规划解读:微架构迭代与ARM服务器落地实践
平头哥倚天720/730/750三代CPU规划解读:微架构迭代与ARM服务器落地实践

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 8:18:31

SQL Explorer 实战指南:在 RocketRide 管道中浏览数据库、编写 SQL 与解读查询计划
SQL Explorer 实战指南:在 RocketRide 管道中浏览数据库、编写 SQL 与解读查询计划

【免费下载链接】rocketride-server High-performance AI pipeline engine with a C core and 50 Python-extensible nodes. Build, debug, and scale LLM workflows with 13 model providers, 8 vector databases, and agent orchestration, all from your IDE. Includes VS C… · 2026/9/25 8:18:31

XAgent utils 模块深度解析:Token 计数、文本裁剪、状态码枚举、任务数据结构与单例元类
XAgent utils 模块深度解析:Token 计数、文本裁剪、状态码枚举、任务数据结构与单例元类

AI Agent大模型后端任务调度 【免费下载链接】XAgent An Autonomous LLM Agent for Complex Task Solving 项目地址: https://gitcode.com/gh_mirrors/xa/XAgent 点击查看 免费下载 XAgent/utils.py 是 XAgent 框架中一个"小而关键"的基础模块&#xff1… · 2026/9/25 8:18:31

数值优化(Numerical Optimization)学习系列-03-共轭梯度方法(Conjugate Gradient)
数值优化(Numerical Optimization)学习系列-03-共轭梯度方法(Conjugate Gradient)

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 1:00:31

创维E900V22D刷机全攻略:S905L3SB芯片兼容性解析与救砖实战
创维E900V22D刷机全攻略:S905L3SB芯片兼容性解析与救砖实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 1:00:31

MQTT协议原理与Broker服务器搭建实战:从Mosquitto到EMQX
MQTT协议原理与Broker服务器搭建实战:从Mosquitto到EMQX

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 1:00:37

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码