ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

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

用 differential dataflow 跑 TPC-H:tpchlike 流式基准评测的工程设计与吞吐实测解读 用 differential dataflow 跑 TPC-Htpchlike 流式基准评测的工程设计与吞吐实测解读【免费下载链接】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 ContestSIGMOD 2016中的思路——那份工作同样把 TPC-H 负载适配到流式场景以模拟真实世界数据不断到来的情形。核心改动在于加载方式传统 TPC-H 先一次性加载全部基础关系再执行查询而流式变体则是把元组源源不断送入各基础关系。README 明确其注入策略round-robin轮询地在各关系之间逐个追加元组即每张表轮流到达一条数据从而营造多路输入同时演进的态势。这一场景与 differential dataflow 的定位天然契合它维护的是数据集合随时间的累积版本每次新元组到来都作为集合上的一次更新diff被吸收后续查询结果随之增量演化——这正是流式数据分析想要的形态。仓库结构一个 crate、八大关系、二十二个查询tpchlike 是一个独立的 Rust crate代码布局如下Cargo.toml声明依赖本地的differential-dataflowpath ../、Git 依赖的timely与arrayvec以及abomonation零拷贝序列化、regex、core_affinity线程绑核等release profile 使用panic abortsrc/types.rsTPC-H 八张基础关系customer、lineitem、nation、orders、part、partsupp、region、supplier的紧凑行类型与.tbl文本解析逻辑src/lib.rs实验公共设施包括InputHandles八路输入句柄、Collections八张表的差分集合包装并记录每张表是否被某查询使用、Arrangements预排序/预索引的 trace与Experimentsrc/queries/q01 至 q22 共 22 个查询模块每个模块内部通常同时给出普通query版本与针对预排序输入的query_arranged版本src/bin/多个可执行入口包括stream.rsREADME 命令行对应的流式评测、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各参数语义结合源码进一步展开pathTPC-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.rsload据此跳过无关文件。logical_batch合并输入的轮数它会改变计算结果。README 建议大部分情况下使用1此时相当于每条元组被独立引入、单独构成一个逻辑时戳。若大于 1多条元组会被合并进同一逻辑批次语义上相当于放宽了到达顺序的粒度。physical_batch一次并发引入的逻辑轮数量。增大它可提升吞吐但以延迟为代价不会改变计算结果——因为差分系统对计算语义无影响只是把更多工作量打包进单次调度。query number1 到 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 Fromstr按 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用分为单位的整数规避浮点误差日期打包为 u32Date是u32用create_date(year, month, day)把年月日分别按 16/8/8 位移进一个整数方便字典序直接比较Order、LineItem的日期字段均如此处理unsafe_abomonate!派生零拷贝序列化借助abomonation宏让行结构可直接在线程/机器间零拷贝传输并在 Cargo.toml 中为ArrayString包装了AbomonationWrapper以适配宏要求RcLineItem共享所有权lineitem是体量最大的表lib.rs 的Collections将其元素类型封装为RcLineItembatch.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 Q1Pricing 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、再叠加税率的*(100tax)/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 10README 公布了一组基于scale factor 10数据集的测量结果——约 10GB 数据lineitem关系含六千万元组——通过变化 physical batching从 1K 元组并发到 1M 元组并发得到。下表同时列出Hot Dog Eating Contest论文中单线程实现的数据注意这些数值仅用于定性对照查询1K1MHot Dog论文单线程query013.76M/s2.67M/s1.27M/squery021.80M/s3.46M/s756.61K/squery033.85M/s8.35M/s3.74M/squery043.00M/s4.47M/s10.08M/squery052.22M/s5.04M/s584.26K/squery0622.77M/s65.23M/s138.33M/squery072.02M/s5.77M/s650.65K/squery081.15M/s2.82M/s91.22K/squery09896.25K/s2.12M/s104.37K/squery103.05M/s9.00M/s2.89M/squery117.02K/s7.90K/s768/squery127.37M/s17.41M/s8.68M/squery13528.03K/s892.74K/s779.52K/squery149.05M/s33.43M/s33.04M/squery151.52M/s6.39M/s17/squery161.82M/s2.91M/s123.94K/squery171.18M/s2.41M/s379.30K/squery183.28M/s4.95M/s1.13M/squery197.48M/s24.61M/s1.95M/squery204.12M/s10.39M/s977/squery21734.42K/s1.39M/s836.80K/squery2221.33K/s11.61K/s189/s数据呈现出两个明显规律physical_batch 从 1K 放大到 1M大多数查询吞吐提升数倍。例如 q06 从 22.77M/s 升到 65.23M/sq19 从 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文本放在某个目录下程序按pathcustomer.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),仅供参考
返回列表