news 2026/9/8 20:34:39

用 differential dataflow 跑 TPC-H:tpchlike 流式基准评测的工程设计与吞吐实测解读

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
用 differential dataflow 跑 TPC-H:tpchlike 流式基准评测的工程设计与吞吐实测解读

用 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-dataflowpath = "../")、Git 依赖的timelyarrayvec,以及abomonation(零拷贝序列化)、regexcore_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.rsstream-concurrent.rsjust-arrange.rssosp.rs等变体,用来对照不同注入与预排序策略。

从 src/bin/stream.rs 的参数解析可以看出,README 中给出的运行命令对应stream这个二进制目标:它跳过前 4 个参数后把剩余参数交给timely::execute_from_args做分布式/多 worker 配置,并将prefixlogical_batchphysical_batchquery分别读为第 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.tbllineitem.tblnation.tblorders.tblpart.tblpartsupp.tblregion.tblsupplier.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]
  • 金额转整数存储acctbalextended_pricediscounttaxsupplycost等小数金额一律(parse::<f64>() * 100.0) as i64,用"分"为单位的整数规避浮点误差;
  • 日期打包为 u32Dateu32,用create_date(year, month, day)把年月日分别按 16/8/8 位移进一个整数,方便字典序直接比较(OrderLineItem的日期字段均如此处理);
  • 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 - intervallineitem行求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都是一组可增量更新的统计量,这里把quantityextended_priceextended_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,正文用explodejoinfiltercount_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论文中单线程实现的数据,注意这些数值仅用于定性对照

查询1K1MHot Dog(论文单线程)
query013.76M/s2.67M/s1.27M/s
query021.80M/s3.46M/s756.61K/s
query033.85M/s8.35M/s3.74M/s
query043.00M/s4.47M/s10.08M/s
query052.22M/s5.04M/s584.26K/s
query0622.77M/s65.23M/s138.33M/s
query072.02M/s5.77M/s650.65K/s
query081.15M/s2.82M/s91.22K/s
query09896.25K/s2.12M/s104.37K/s
query103.05M/s9.00M/s2.89M/s
query117.02K/s7.90K/s768/s
query127.37M/s17.41M/s8.68M/s
query13528.03K/s892.74K/s779.52K/s
query149.05M/s33.43M/s33.04M/s
query151.52M/s6.39M/s17/s
query161.82M/s2.91M/s123.94K/s
query171.18M/s2.41M/s379.30K/s
query183.28M/s4.95M/s1.13M/s
query197.48M/s24.61M/s1.95M/s
query204.12M/s10.39M/s977/s
query21734.42K/s1.39M/s836.80K/s
query2221.33K/s11.61K/s189/s

数据呈现出两个明显规律:

  1. physical_batch 从 1K 放大到 1M,大多数查询吞吐提升数倍。例如 q06 从 22.77M/s 升到 65.23M/s,q19 从 7.48M/s 升到 24.61M/s。这正对应 README 对参数的解释:更大的物理批把更多轮次合并并发,摊薄了调度与同步开销,但会牺牲逐元组的低延迟。少数查询(如 q01、q22)反而出现回落,说明批大小与具体算子的增量维护成本存在复杂的权衡,不能一概而论。
  2. 横向对比论文单线程结果可见"改进空间的位置":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 表达能力与增量更新机制的工作台,而非已经校准过的基准工具。

复现实验的快速路径

如果你想在本地复现或扩展这套评测,可以按以下步骤进行(仓库为只读,仅用于查看与运行):

  1. 阅读 external/differential-dataflow/tpchlike/README.md 及本指南,确认实验意图与参数含义;
  2. 准备 TPC-H 数据:按 scale factor 用 dbgen 生成八张.tbl文本,放在某个目录下(程序按<path>customer.tbl的方式拼接文件名,注意结尾需自带路径分隔符);
  3. 以 release 模式运行流式评测,例如查询 q06 且按元组独立引入:
cd external/differential-dataflow/tpchlike cargo run --release -- /path/to/tpch-sf10/ 1 1000000 6
  1. 对照阅读 src/bin/stream.rs 理解注入循环,必要时修改physical_batch观察吞吐-延迟曲线;使用 src/bin/batch.rs 可获得全量批装载的对照读数;
  2. 若只关心某条查询的逻辑,可直接查看 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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/8 20:26:59

Claude Code 十大实战技能:从安装配置到Skills定制与Token成本控制

Claude Code 这个终端里的 AI 编程 Agent&#xff0c;最近几乎把所有做开发的朋友都圈进来了。它跟 IDE 里那些只做代码补全的插件完全不同&#xff0c;是一个能真正看懂整个项目结构、自己动手改文件、跑测试、提交 Git 的智能体。过去大半年我把这个工具从安装、配置到深度定…

作者头像 李华
网站建设 2026/9/8 20:25:49

零改板替换实战:VL171换国产CSA171的踩坑全记录与实操指南

零改板替换&#xff0c;我把VL171换成了国产CSA171&#xff1a;踩坑全记录与实操指南国产芯片替代这个话题&#xff0c;这两年在硬件圈里几乎天天有人在聊。但我发现一个现象&#xff1a;很多人一提到“零改板替换”&#xff0c;第一反应就是“引脚对得上就行”&#xff0c;结果…

作者头像 李华
网站建设 2026/9/8 20:25:01

用 tiny11builder 精简 Windows 11 安装镜像:ISO 体积缩减 41.8%

用 tiny11builder 精简 Windows 11 安装镜像&#xff1a;ISO 体积缩减 41.8% 【免费下载链接】tiny11builder Scripts to build a trimmed-down Windows 11 image. 项目地址: https://gitcode.com/GitHub_Trending/ti/tiny11builder tiny11builder 是一个纯 PowerShell …

作者头像 李华
网站建设 2026/9/8 20:24:42

条件GAN在垃圾邮件数据填补中的实战应用

简介&#xff1a;本资源是一份面向深度学习初学者与数据科学实践者的GAN缺失值填补实战代码包&#xff0c;聚焦Spam邮件数据集中的缺失特征修复问题&#xff0c;适用于机器学习预处理、学术研究及课程设计等场景。压缩包共2个文件&#xff08;127KB&#xff09;&#xff0c;含核…

作者头像 李华
网站建设 2026/9/8 20:24:20

十分钟完成第一次捕获:res-downloader 资源下载器上手教程

十分钟完成第一次捕获&#xff1a;res-downloader 资源下载器上手教程 【免费下载链接】res-downloader 视频号、小程序、抖音、快手、小红书、直播流、m3u8、酷狗、QQ音乐等常见网络资源下载! 项目地址: https://gitcode.com/GitHub_Trending/re/res-downloader 想把一…

作者头像 李华