ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

RisingWave 扩展查询模式 E2E 测试:参数绑定、MaxRows 与查询取消的实战验证

RisingWave 扩展查询模式 E2E 测试:参数绑定、MaxRows 与查询取消的实战验证 数据库流处理后端数据工程【免费下载链接】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 是一款面向实时数据流的流式处理数据库其查询能力既包含常规的简单查询协议也支持 PostgreSQL 风格的扩展查询协议Extended Query Protocol。本文以仓库中 src/tests/e2e_extended_mode/README.md 为骨架深入讲解 RisingWave 为覆盖 sqllogictest 无法触达的扩展查询模式场景而专门构建的端到端测试程序risingwave_e2e_extended_mode_test它验证了什么、为什么需要独立于 sqllogictest 存在、如何编译与运行以及每个测试用例背后的协议原理与源码实现。读完本文你将理解 RisingWave 扩展查询模式下的参数绑定Bind、最大返回行数MaxRows、查询取消Cancel以及订阅游标Subscription Cursor取数的完整测试思路并能在本地复现这套验证流程。为什么需要独立的扩展模式 E2E 测试程序RisingWave 的日常端到端测试大量依赖 sqllogictest例如 e2e_test/extended_mode/README.md 中给出的命令就是用 sqllogictest 以扩展查询模式执行.slt用例sqllogictest -p 4566 -d dev -e postgres-extended ./e2e_test/extended_mode/**/*.slt然而扩展查询模式Extended Query Protocol下有一批关键行为是 sqllogictest 这种以结果集文本比对为主的工具无法覆盖的。原文档明确指出以下三类能力无法在 sqllogictest 中测试bind parameter绑定参数客户端在 Parse 阶段将语句中的$1、$2等占位符与具体值绑定sqllogictest 的 DSL 难以表达这种先解析、再绑定、后执行的过程化交互。max row number最大返回行数扩展协议中 Execute 消息可携带MaxRows用于限制单次执行返回的行数配合游标Portal实现分批取数sqllogictest 没有对应的行数上限控制原语。cancel query取消查询通过CancelRequest消息异步中止正在执行的查询属于纯过程化、时序敏感的交互基于文本回放的 sqllogictest 无法稳定复现。原文档引用 PostgreSQL 15 的协议文档指出扩展协议的精髓在于 Portal 与 MaxRows 的组合语义——一旦 Portal 被创建后续的 Execute 可以反复向它发送不同 MaxRows 的取数请求而返回行数计数互不干扰。这段语义正是 RisingWave 需要独立程序来专门验证的原因。因此仓库在 src/tests/e2e_extended_mode 下维护了一个独立的 Rust 二进制测试程序。原文档也明确说明了它的定位在 sqllogictest 支持这些能力之前由本程序承担这些函数的测试未来这些能力有可能被合并回 e2e_test/extended_mode 目录下的 sqllogictest 用例中。程序结构一个轻量的 tokio tokio-postgres 测试套件整个测试程序只有 5 个源文件规模克制、职责清晰src/tests/e2e_extended_mode/src/main.rs程序入口解析命令行参数后调用run_test_suit以进程退出码反映测试成败src/tests/e2e_extended_mode/src/opts.rs基于 clap 定义命令行参数src/tests/e2e_extended_mode/src/lib.rsrun_test_suit的封装负责初始化日志、驱动测试套件并汇总结果src/tests/e2e_extended_mode/src/test.rs核心测试逻辑全部 11 组用例都定义在这里src/tests/e2e_extended_mode/Cargo.toml依赖清单。从依赖可以看到这套测试的底层技术选型tokio提供异步运行时与tokio::spawn并发执行查询tokio-postgres作为 PostgreSQL 线协议客户端Client、NoTlsclap负责命令行解析chrono / pg_interval / rust_decimal提供日期时间、区间与十进制数值的 Rust 类型映射tracing / tracing-subscriber输出测试日志。其中tokio-postgres的features [with-chrono-0_4]表明测试会通过 chrono 类型直接读写数据库的日期时间列这在后面的二进制参数用例中会实际用到。入口函数 src/tests/e2e_extended_mode/src/main.rs 的完整逻辑是用Opts::parse()解析参数初始化 tracing 日志然后依次把db_name / user_name / host / port / password传给run_test_suit最后用exit(...)把测试结果转换为进程退出码0 成功、1 失败。而 src/tests/e2e_extended_mode/src/lib.rs 中的run_test_suit在失败时会打印提示Risingwave e2e extended mode test failed ... Please ensure that your psql version is larger than 14.1暗示这套测试对客户端协议特性有一定版本要求psql 14.1 之后才具备更完整的扩展协议交互能力。命令行参数详解与运行方式原文档给出的运行命令是RUST_BACKTRACE1 target/debug/risingwave_e2e_extended_mode_test --host 127.0.0.1 \ -p 4566 \ -u root \ --database dev \对照 src/tests/e2e_extended_mode/src/opts.rs每个参数的完整语义如下表命令行参数属性名源码默认值说明--host PG_SERVER_ADDRESSpg_server_hostlocalhost待测 RisingWave 服务地址-p, --port PG_SERVER_PORTpg_server_port无必填待测服务端口如 4566-u, --user PG_USERNAMEpg_user_namepostgres连接用户名原文档示例中显式使用root--database DBpg_db_namedev连接使用的数据库名--password PG_PASSWARDpg_password空字符串连接密码为空时不写入连接串关于密码有一个值得注意的实现细节在 src/tests/e2e_extended_mode/src/test.rs 中TestSuite::new会根据密码是否为空决定连接串格式——有密码时拼接dbname... user... host... port... password...无密码时省略password字段。也就是说本地无鉴权环境可以直接省略--password。在 CI 中该程序的实际调用方式见 ci/scripts/e2e-test-serial.sh与文档命令几乎一致只是省略了--database使用默认的devecho --- e2e, $mode, extended query cluster_start risedev slt -p 4566 -d dev -e postgres-extended ./e2e_test/extended_mode/**/*.slt RUST_BACKTRACE1 target/debug/risingwave_e2e_extended_mode_test --host 127.0.0.1 \ -p 4566 \ -u root可见在 CI 流程中sqllogictest 的扩展模式用例与这个独立程序是先后串联执行的先跑 e2e_test/extended_mode 下的文本用例再跑本程序覆盖协议级行为。构建方面ci/scripts/build.sh 将risingwave_e2e_extended_mode_test作为 CI 产出的核心工件之一与其他测试工具一同编译并上传说明它是正式回归链路的一部分。测试套件的整体编排src/tests/e2e_extended_mode/src/test.rs 中TestSuite::test()按顺序串联了 11 组用例构成一条完整回归链binary_param_and_result—— 二进制格式参数绑定与结果回读dql_dml_with_param—— 带参数的 DQL/DMLINSERT/UPDATE/DELETE/SELECTmax_row—— 事务内 Portal 的 MaxRows 分批取数multiple_on_going_portal—— 同一事务内多个并发 Portal 独立推进create_with_parameter—— 验证CREATE TABLE/VIEW AS SELECT $1应报错负向用例simple_cancel(false)/simple_cancel(true)—— 简单查询取消local / distributed 两种 query_modecomplex_cancel(false)/complex_cancel(true)—— 复杂多表 JOIN 查询取消subscription_fetch_cancel(false)/subscription_fetch_cancel(true)—— 订阅游标取数取消subquery_with_param—— 子查询内的参数绑定create_mview_with_parameter—— 带参数创建物化视图及其重命名。其中query_mode是 RisingWave 特有的会话参数local为单机本地执行、distributed为分布式执行测试通过 create_client 在连接建立后执行set query_mode local/distributed来切换两种执行路径从而保证协议特性在两种执行模式下都成立。测试中还大量使用自制的test_eq!宏见 src/tests/e2e_extended_mode/src/test.rs做逐值断言失败时输出文件与行号便于定位。二进制参数绑定与类型往返验证扩展查询模式最基础的能力是先 Prepare 再 Bind客户端把 SQL 中的$1、$2等占位符与类型化参数值绑定服务端执行后再把结果以二进制编码返回。binary_param_and_result用例逐一验证了 RisingWave 与 tokio-postgres 之间的类型往返正确性src/tests/e2e_extended_mode/src/test.rs测试 SQL绑定参数类型期望回读值select $1::SMALLINTi161024select $1::INTi32144232select $1::BIGINTi6499999999select $1::DECIMALrust_decimal::Decimal2.33454select $1::BOOLbooltrueselect $1::REALf321.234234select $1::DOUBLE PRECISIONf64234234.23490238483select $1::dateNaiveDate2022-01-01select $1::timeNaiveTime10:00:00select $1::timestampNaiveDateTime2022-01-01 10:00:00select $1::timestamptzDateTimeUtc2022-01-01 10:00:00 (UTC)select $1::intervalpg_interval::Interval1 年 1 月 24000000 微秒这里的重点是query()传入的参数数组[value]会触发 tokio-postgres 走扩展协议中的 Parse → Bind → Execute 全流程而try_get::_, T则要求服务端以对应类型的二进制编码返回结果。因此这组用例同时验证了两件事参数值能被 RisingWave 正确解析并参与表达式计算以及返回结果的二进制编解码与客户端类型完全一致。interval用例通过pg_interval::Interval::new(1, 1, 24000000)构造 1 年 1 月加 2400 万微秒的区间值覆盖了复合时间单位的编解码。带参数的 DQL/DMLPrepare 语句的生命周期dql_dml_with_param用例src/tests/e2e_extended_mode/src/test.rs验证参数化语句在数据操作链路中的完整生命周期创建表t(id int)用prepare_typed(insert INTO t (id) VALUES ($1), [])准备插入语句循环绑定0..20共 20 行再执行flush强制落盘prepare_typed(select * FROM t where id $1 order by id ASC, [Type::INT4])带类型准备查询绑定10后应返回id 0..9共 10 行prepare_typed(update t set id $1 where id $2, [Type::INT4, Type::INT4])绑定(100, 10)更新全部 10 行后同样条件再查应返回 0 行prepare_typed(delete FROM t where id $1, [Type::INT4])绑定20删除所有行再以101为阈值查询此时id 100的行应保留 10 行drop table t收尾。值得注意的细节是这里prepare_typed显式传入Type::INT4指定占位符类型而 INSERT 语句则传空类型列表让服务端推断。用例同时覆盖了扩展协议中 PrepareStatement对象可复用与 Execute同一语句多次绑定执行的语义且每次数据变更后都调用flush确保流式处理引擎的数据已持久化可见避免读到过期快照。MaxRows 与 Portal分批取数的正确语义max_row用例src/tests/e2e_extended_mode/src/test.rs是对原文档所引 PostgreSQL 协议语义的直接实现。它在一个事务中演示了扩展协议最独特的能力——同一条 Statement 绑定出一个 Portal多次 Execute 用不同的 MaxRows 逐批取出数据let transaction client.transaction().await?; let statement transaction .prepare_typed(SELECT * FROM t order by id, []) .await?; let portal transaction.bind(statement, []).await?; // 前 5 次每次 MaxRows1逐行取出 id 0..4 for t in 0..5 { let rows transaction.query_portal(portal, 1).await?; test_eq!(rows.len(), 1); let id: i32 row.get(0); test_eq!(id, t); } // 然后 MaxRows3一次性取 3 行id 5,6,7 let mut i 5; for row in transaction.query_portal(portal, 3).await? { let id: i32 row.get(0); test_eq!(id, i); i 1; } test_eq!(i, 8); // 最后 MaxRows5取剩余 2 行id 8,9Portal 耗尽 for row in transaction.query_portal(portal, 5).await? { ... } test_eq!(i, 10);这正好对应 PostgreSQL 文档中Execute 消息里的 MaxRows 只影响本次取数Portal 的返回行数计数独立于新一次的 MaxRows而且取完即终止的行为。若 RisingWave 对 MaxRows 的实现有偏差例如一次取多、忽略上限、或计数被重置上述断言会立刻失败。多 Portal 并发推进状态隔离验证multiple_on_going_portal用例src/tests/e2e_extended_mode/src/test.rs验证同一个事务中两个 Portal 各自独立推进、互不干扰同一个 StatementSELECT generate_series(1,5,1)绑定出portal_1和portal_2对portal_1取 1 行 → 得到 1对portal_2取 1 行 → 得到 1portal_2 从自己的位置开始而不是接着 portal_1对portal_2取 3 行 → 得到 2、3、4再对portal_1取 1 行 → 得到 2portal_1 仍然在自己上次取数后的位置。这个用例确保 RisingWave 为每个 Portal 维护独立的执行位置这是扩展协议中多游标并发例如同一事务内交替消费多个结果集正确性的基石。参数化 DDL 的限制与物化视图参数的例外负向用例CREATE TABLE/VIEW AS SELECT $1 必须报错create_with_parameter用例src/tests/e2e_extended_mode/src/test.rs是一组负向测试注释明确写着 Cant support these sqltest_eq!(client.query(create table t as select $1, []).await.is_err(), true); test_eq!(client.query(create view v as select $1, []).await.is_err(), true);即扩展协议下CREATE TABLE AS SELECT $1与CREATE VIEW AS SELECT $1这样的参数化 DDL 在 RisingWave 中是不支持的应当返回错误。这类断言锁定了当前行为边界防止未来改动意外放开。正向例外CREATE MATERIALIZED VIEW 支持参数与上述限制形成对照的是create_mview_with_parameter用例src/tests/e2e_extended_mode/src/test.rs验证物化视图创建语句可以带参数let statement client .prepare_typed(create materialized view mv as select $1 as x, [Type::INT4]) .await?; client.execute(statement, [42_i32]).await?; // select * from mv → 1 行值为 42注释揭示了实现原理Test renaming mv because it relies on parsing and rewrite thecreate MVquery——RisingWave 对CREATE MATERIALIZED VIEW ... AS SELECT $1的处理是先解析并改写创建语句因此参数能被正确内联。用例还顺带验证了改写逻辑的健壮性执行ALTER MATERIALIZED VIEW mv RENAME TO mv2后从mv2查询仍能得到 42说明重命名没有破坏参数化创建的产物。查询取消simple / complex 两种难度简单查询取消simple_cancel(is_distributed)用例src/tests/e2e_extended_mode/src/test.rs流程创建表并插入 1000 行数据从client.cancel_token()取得取消令牌tokio::spawn在一个独立任务中发起select * from t全表扫描足够慢用tokio::select!竞争两个分支查询完成说明取消失败与cancel_token.cancel_query(NoTls)成功说明取消成功用新连接验证取消后的数据一致性select * from t order by id limit 10应返回 0..9证明取消没有破坏表数据或连接。注意simple_cancel会以falselocal 模式和truedistributed 模式各跑一遍覆盖两种执行引擎。复杂查询取消complex_cancel用例src/tests/e2e_extended_mode/src/test.rs把难度提升到多表 JOINt1、t2、t3 三张表各插入 1000 行然后执行一个包含 INNER JOIN、子查询、LIKE 过滤、GROUP BY MAX、LEFT JOIN、ORDER BY 的复杂 SQL。取消成功后用新连接对同一 SQL 加LIMIT 10重放断言返回前 10 行恰好是(1,1,1)、(10,10,10)、(100,100,100)、(101,101,101) …。这一组用例验证的是复杂执行计划构建过程中或执行中收到 CancelRequest 时能干净地终止且不影响后续同类查询的正确结果。订阅游标取数取消流式场景下的取消语义subscription_fetch_cancel用例src/tests/e2e_extended_mode/src/test.rs把取消测试延伸到 RisingWave 特有的**订阅Subscription**功能这是普通关系型数据库没有的场景建表sub_cancel_t_local/dist创建订阅create subscription ... with(retention 1D)declare cur subscription cursor for sub since now()声明订阅游标后台任务执行fetch 1 from cur with (timeout 60s)阻塞等待流数据1 秒后调用cancel_token.cancel_query(NoTls)用 10 秒超时包裹 fetch 任务断言它必须以错误结束——若 fetch 成功返回则用例报错 subscription cursor fetch should be cancelled by CancelRequest清理drop subscription、drop table。该用例同样对 local / distributed 两种 query_mode 各跑一遍。它验证的核心语义是阻塞在流式游标取数上的查询也能被CancelRequest中断这是流处理引擎在扩展协议交互上的关键正确性保证。子查询中的参数绑定与回归链路总结subquery_with_paramsrc/tests/e2e_extended_mode/src/test.rs是最简用例select (select $1::SMALLINT)在子查询中绑定i16参数 1024 并回读验证参数解析在嵌套子查询中依然生效。将 11 组用例放在一起可以看到这套程序覆盖了扩展查询协议在 RisingWave 中的完整行为面参数绑定标量类型往返、DML 语句复用、子查询内参数、物化视图参数正向/负向双面验证Portal 与 MaxRows分批取数语义、多 Portal 状态隔离查询取消简单查询、复杂 JOIN、流式订阅游标且全部覆盖 local 与 distributed 两种执行模式数据一致性每次取消后都用新连接复验结果与表数据。如何在本地复现这套测试要完整跑通上述验证需要依次完成构建程序在 RisingWave 仓库根目录执行cargo build -p risingwave_e2e_extended_mode_testCI 中的完整构建可参考 ci/scripts/build.sh会附加--features udf --features datafusion等特性。产物为target/debug/risingwave_e2e_extended_mode_test。启动 RisingWave用仓库自带的risedev工具启动本地集群如risedev dev确保 4566 端口可达。运行程序按 src/tests/e2e_extended_mode/README.md 的命令执行RUST_BACKTRACE1 target/debug/risingwave_e2e_extended_mode_test --host 127.0.0.1 \ -p 4566 \ -u root \ --database dev \对照检查若全部用例通过日志会输出 Risingwave e2e extended mode test completed successfully! 且退出码为 0若失败src/tests/e2e_extended_mode/src/lib.rs 会打印错误详情与提示例如客户端版本需大于 14.1退出码为 1。另外正如 CI 脚本 ci/scripts/e2e-test-serial.sh 所示sqllogictest 扩展模式用例与本程序是配套执行的关系建议先跑risedev slt -p 4566 -d dev -e postgres-extended ./e2e_test/extended_mode/**/*.slt覆盖文本层用例再跑本程序覆盖协议层行为形成完整闭环。赞分享数据库流处理后端数据工程【免费下载链接】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 扩展查询模式Extended Query Mode端到端测试指南RisingWave 扩展查询模式Extended Query Mode端到端测试指南 本指南聚焦 RisingWave 的 PostgreSQL 扩展查询数据库流处理后端数据工程RisingWave 联邦查询实战通过 Presto/Trino PostgreSQL 连接器查询 RisingWave 数据RisingWave 联邦查询实战通过 Presto/Trino PostgreSQL 连接器查询 RisingWave 数据 本文基于 RisingWave数据库流处理后端数据工程MobileOne_s4.apple_in1k终极1毫秒移动骨干网络如何实现ImageNet-1k高效图像分类MobileOne_s4.apple_in1k终极1毫秒移动骨干网络如何实现ImageNet 1k高效图像分类 在移动设备上实现快速高效的图像分类一直是计上一篇Coil性能巅峰终极内存优化与磁盘缓存策略全解析下一篇分布式采样策略ESFT多GPU数据加载优化创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表