用 differential dataflow 跑 TPC-H:tpchlike 流式基准评测的工程设计与吞吐实测解读
【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway
本篇技术指南聚焦于仓库external/differential-dataflow/tpchlike子项目,它把经典数据库基准 TPC-H 改造成流式增量场景,用于评测 differential dataflow(差分数据流)在持续追加数据时的计算性能。读完本文,你将掌握该基准的实验动机、cargo run --release -- <path> <logical_batch> <physical_batch> <query number>的完整参数语义、22 个查询的实现组织方式,以及如何在 scale factor 10(约 10GB、lineitem六千万元组)数据集上复现并解读吞吐量测量结果。
为什么要把 TPC-H 改造成"流式"评测
tpchlike 的 README 开宗明义:TPC-H 是被"严肃人士"用来评估"能赚钱的系统"的数据库基准,而 differential dataflow 虽然没那么严肃,人们仍然好奇它在这类任务上的表现。其评测设计模仿了论文How to Win a Hot Dog Eating Contest(SIGMOD 2016)中的思路——那份工作同样把 TPC-H 负载适配到流式场景,以模拟真实世界数据不断到来的情形。
核心改动在于加载方式:传统 TPC-H 先一次性加载全部基础关系再执行查询;而流式变体则是把元组源源不断送入各基础关系。README 明确其注入策略:round-robin(轮询)地在各关系之间逐个追加元组,即每张表轮流到达一条数据,从而营造多路输入同时演进的态势。
这一场景与 differential dataflow 的定位天然契合:它维护的是"数据集合随时间的累积版本",每次新元组到来都作为集合上的一次更新(diff)被吸收,后续查询结果随之增量演化——这正是流式数据分析想要的形态。
仓库结构:一个 crate、八大关系、二十二个查询
tpchlike 是一个独立的 Rust crate,代码布局如下:
- Cargo.toml:声明依赖本地的
differential-dataflow(path = "../")、Git 依赖的timely与arrayvec,以及abomonation(零拷贝序列化)、regex、core_affinity(线程绑核)等;release profile 使用panic = "abort"; - src/types.rs:TPC-H 八张基础关系(customer、lineitem、nation、orders、part、partsupp、region、supplier)的紧凑行类型与
.tbl文本解析逻辑; - src/lib.rs:实验公共设施,包括
InputHandles(八路输入句柄)、Collections(八张表的差分集合包装,并记录每张表是否被某查询使用)、Arrangements(预排序/预索引的 trace)与Experiment; - src/queries/:q01 至 q22 共 22 个查询模块,每个模块内部通常同时给出普通
query版本与针对预排序输入的query_arranged版本; - src/bin/:多个可执行入口,包括
stream.rs(README 命令行对应的流式评测)、batch.rs(一次性批量注入的对比实现)以及arrange.rs、stream-concurrent.rs、just-arrange.rs、sosp.rs等变体,用来对照不同注入与预排序策略。
从 src/bin/stream.rs 的参数解析可以看出,README 中给出的运行命令对应stream这个二进制目标:它跳过前 4 个参数后把剩余参数交给timely::execute_from_args做分布式/多 worker 配置,并将prefix、logical_batch、physical_batch、query分别读为第 1、2、3、4 个命令行参数;此外还可选追加seal-inputs标记,用于让系统在处理完数据后关闭输入。
命令行参数:path、logical_batch、physical_batch 与查询编号
README 规定的运行方式如下:
cargo run --release -- <path> <logical_batch> <physical_batch> <query number>各参数语义(结合源码进一步展开):
<path>:TPC-H 数据文件所在目录前缀。程序会在该目录下寻找customer.tbl、lineitem.tbl、nation.tbl、orders.tbl、part.tbl、partsupp.tbl、region.tbl、supplier.tbl这些标准|分隔的文本表。README 提示:如果手头没有这些文件,可以用 TPC-H 官方发布的 dbgen 生成器按 scale factor 产出。实际加载只发生在当前查询真正用到的那几张表上——Collections通过used: [bool; 8]记录每张表是否被访问(见 src/lib.rs),load据此跳过无关文件。<logical_batch>:合并输入的"轮数",它会改变计算结果。README 建议大部分情况下使用1,此时相当于每条元组被独立引入、单独构成一个逻辑时戳。若大于 1,多条元组会被合并进同一逻辑批次,语义上相当于放宽了到达顺序的粒度。<physical_batch>:一次并发引入的逻辑轮数量。增大它可提升吞吐但以延迟为代价,不会改变计算结果——因为差分系统对计算语义无影响,只是把更多工作量打包进单次调度。<query number>:1 到 22 的整数,选择执行哪个 TPC-H 查询。
真正驱动这些批次的逻辑藏在 src/bin/stream.rs 的load函数中。它对每一行文本计算:
let logical = (8 * count / logical_batch) + off; // off 为各表序号 0..7 let physical = logical / physical_batch; let round = physical / 8;off是当前表在八张关系中的序号,配合系数 8,正好实现了 README 所说的"round-robin 逐元组轮询注入":行号为count的元组被分配到逻辑轮8*count / logical_batch,随后物理批physical_batch又决定多少逻辑轮合并为一次并发注入。主循环中每个物理批先send_batch注入各表剩余数据,再统一advance_to到下一个轮次并worker.step_while推进到对应时间,直到全部数据消费完毕(src/bin/stream.rs)。
评测的计时与计量方式也同样体现在源码中:所有表完成加载并推进到同一轮次后启动Instant::now()计时器,最后按rate = (peers * tuples) / (elapsed_seconds)计算吞吐,并在 0 号 worker 上打印一行制表符分隔的结果:query_name、logical_batch、physical_batch、peers、rate、耗时(ns)(src/bin/stream.rs)。也就是说,每条数据按"表内被当前 worker 处理过的元组总量"计量,从而可以直接给出元组/秒的量纲。
与之相对的 src/bin/batch.rs 只接收<path> <query>两个参数,把每张表整体作为一个物理批一次性send_batch注入,对应经典 TPC-H 的"先全量装载再查询"模式,可作为流式方案(stream)与批式装载之间的对照实验。
从.tbl文本到紧凑内存行:types.rs 的数据工程
为了让八张表的文本能高效灌入数据流,types.rs 做了一系列值得借鉴的取舍:
|分隔解析:每个行类型实现impl From<&str>,按 TPC-H 标准的|分隔字段逐个split解析;- 紧凑定长存储:字段尽量使用
[u8; N]定长数组或arrayvec::ArrayString固定容量字符串,避免堆分配,例如[u8; 25]存厂商名、[u8; 15]存电话;Part甚至把brand存成[u8; 10]; - 金额转整数存储:
acctbal、extended_price、discount、tax、supplycost等小数金额一律(parse::<f64>() * 100.0) as i64,用"分"为单位的整数规避浮点误差; - 日期打包为 u32:
Date是u32,用create_date(year, month, day)把年月日分别按 16/8/8 位移进一个整数,方便字典序直接比较(Order、LineItem的日期字段均如此处理); unsafe_abomonate!派生零拷贝序列化:借助abomonation宏让行结构可直接在线程/机器间零拷贝传输,并在 Cargo.toml 中为ArrayString包装了AbomonationWrapper以适配宏要求;Rc<LineItem>共享所有权:lineitem是体量最大的表,lib.rs 的Collections将其元素类型封装为Rc<LineItem>,batch.rs/stream.rs在注入时用line.map(|(d,t,r)| (Rc::new(d),t,r))包一层引用计数(如 src/bin/stream.rs),让多条派生数据可以共享同一行内容而无需复制。
正是这种"行即定长紧凑结构体、文本解析一次完成"的设计,保证了后续测量反映的是计算本身的开销,而非 I/O 解析的噪声。
一条 SQL 如何变成差分数据流:以 q01 为例
query01.rs 的文件头完整保留了 TPC-H Q1(Pricing Summary Report)的原始 SQL:按l_returnflag, l_linestatus分组,对ship_date <= 1998-12-01 - interval的lineitem行求sum(l_quantity)、sum(l_extendedprice)、带折扣的sum_disc_price、带折扣和税率的sum_charge以及计数。
其差分数据流实现几乎是一行行"翻译":
collections .lineitems() .explode(|item| if item.ship_date <= ::types::create_date(1998, 9, 2) { Some(((item.return_flag[0], item.line_status[0]), DiffPair::new(item.quantity as isize, DiffPair::new(item.extended_price as isize, DiffPair::new((item.extended_price * (100 - item.discount) / 100) as isize, DiffPair::new((item.extended_price * (100 - item.discount) * (100 + item.tax) / 10000) as isize, DiffPair::new(item.discount as isize, 1))))))) } else { None } ) .count_total() .probe_with(probe);要点在于差分库特有的DiffPair嵌套:差分数据流的"值"本身可以携带多重累加字段,每层DiffPair都是一组可增量更新的统计量,这里把quantity、extended_price、extended_price*(1-discount)(整数化后写成price*(100-discount)/100)、再叠加税率的*(100+tax)/10000以及discount层层嵌套,最后经count_total()按(return_flag, line_status)分组完成 SQL 的聚合。同一个query函数还给出了面向预排序输入(query_arranged)的等价变体。其余查询文件(q02 至 q22)遵循同一模式:头部注释保留原 TPC-H SQL,正文用explode、join、filter、count_total等算子重新表达。
查询的底层支撑来自 lib.rs 中的Arrangements:程序启动时对 customer、nation、order、part、partsupp、region、supplier 等表按主键做arrange_by_key()预排序,把 trace 持久化为可被查询import_core复用的索引,并且可以按实验需要set_physical_compaction/set_logical_compaction控制历史版本的压缩策略——这也是同一查询能排出query/query_arranged两种写法的原因,用于验证"把 join 键预先安排成索引"对吞吐的影响。
吞吐量测量与解读(scale factor 10)
README 公布了一组基于scale factor 10数据集的测量结果——约 10GB 数据,lineitem关系含六千万元组——通过变化 physical batching(从 1K 元组并发到 1M 元组并发)得到。下表同时列出Hot Dog Eating Contest论文中单线程实现的数据,注意这些数值仅用于定性对照:
| 查询 | 1K | 1M | Hot Dog(论文单线程) |
|---|---|---|---|
| query01 | 3.76M/s | 2.67M/s | 1.27M/s |
| query02 | 1.80M/s | 3.46M/s | 756.61K/s |
| query03 | 3.85M/s | 8.35M/s | 3.74M/s |
| query04 | 3.00M/s | 4.47M/s | 10.08M/s |
| query05 | 2.22M/s | 5.04M/s | 584.26K/s |
| query06 | 22.77M/s | 65.23M/s | 138.33M/s |
| query07 | 2.02M/s | 5.77M/s | 650.65K/s |
| query08 | 1.15M/s | 2.82M/s | 91.22K/s |
| query09 | 896.25K/s | 2.12M/s | 104.37K/s |
| query10 | 3.05M/s | 9.00M/s | 2.89M/s |
| query11 | 7.02K/s | 7.90K/s | 768/s |
| query12 | 7.37M/s | 17.41M/s | 8.68M/s |
| query13 | 528.03K/s | 892.74K/s | 779.52K/s |
| query14 | 9.05M/s | 33.43M/s | 33.04M/s |
| query15 | 1.52M/s | 6.39M/s | 17/s |
| query16 | 1.82M/s | 2.91M/s | 123.94K/s |
| query17 | 1.18M/s | 2.41M/s | 379.30K/s |
| query18 | 3.28M/s | 4.95M/s | 1.13M/s |
| query19 | 7.48M/s | 24.61M/s | 1.95M/s |
| query20 | 4.12M/s | 10.39M/s | 977/s |
| query21 | 734.42K/s | 1.39M/s | 836.80K/s |
| query22 | 21.33K/s | 11.61K/s | 189/s |
数据呈现出两个明显规律:
- physical_batch 从 1K 放大到 1M,大多数查询吞吐提升数倍。例如 q06 从 22.77M/s 升到 65.23M/s,q19 从 7.48M/s 升到 24.61M/s。这正对应 README 对参数的解释:更大的物理批把更多轮次合并并发,摊薄了调度与同步开销,但会牺牲逐元组的低延迟。少数查询(如 q01、q22)反而出现回落,说明批大小与具体算子的增量维护成本存在复杂的权衡,不能一概而论。
- 横向对比论文单线程结果可见"改进空间的位置":q15、q19、q20、q22 等查询相比论文结果提升明显(q20 从 977/s 提升到 10.39M/s 量级);而 q04、q06 上差分数据流仍不及论文单线程实现,作者明确将其标注为"还有提升空间"的方向。
最重要的使用前提:这份测量不等于"正确答案"
README 末尾以醒目语气给出了必须遵守的免责声明:
这些时间是仓库中代码的运行时间,但代码可能并没有计算出本应计算的量。作者直言自己很可能搞砸了部分甚至全部查询实现——目前没有可对照的基准结果来做正确性验证,因此这些测量值不能被当作事实使用。
这意味着,若要以本仓库的数据作为任何结论的依据,合理的做法是先独立验证各查询输出:可以拿官方 dbgen 的小规模数据(如 SF 1)生成标准 TPC-H 参考结果,与tpchlike各查询模块的输出做逐值比对;也可对照 types.rs 中的字段映射,逐一确认日期比较、金额整数化(乘 100)、折扣/税率公式(整数化后的*(100-discount)/100)与原始 SQL 语义完全一致。作者在 README 中表示欢迎指出 bug 或协助验证——本质上,tpchlike 更应被当作演示 differential dataflow 表达能力与增量更新机制的工作台,而非已经校准过的基准工具。
复现实验的快速路径
如果你想在本地复现或扩展这套评测,可以按以下步骤进行(仓库为只读,仅用于查看与运行):
- 阅读 external/differential-dataflow/tpchlike/README.md 及本指南,确认实验意图与参数含义;
- 准备 TPC-H 数据:按 scale factor 用 dbgen 生成八张
.tbl文本,放在某个目录下(程序按<path>customer.tbl的方式拼接文件名,注意结尾需自带路径分隔符); - 以 release 模式运行流式评测,例如查询 q06 且按元组独立引入:
cd external/differential-dataflow/tpchlike cargo run --release -- /path/to/tpch-sf10/ 1 1000000 6- 对照阅读 src/bin/stream.rs 理解注入循环,必要时修改
physical_batch观察吞吐-延迟曲线;使用 src/bin/batch.rs 可获得全量批装载的对照读数; - 若只关心某条查询的逻辑,可直接查看 src/queries/ 下对应编号文件——每个文件的 SQL 注释即是最好的说明书。
通过这样的流程,你既能快速上手 differential dataflow 在复杂多表 join 聚合负载上的编程范式,也能建立起"流式批大小 ↔ 吞吐"的直观量化认知,为在真实流式分析场景中选型 batch 参数与预排序策略提供参照。
【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考