@Greptime: DataFusion 在去年九月向上游合并了动态过滤下推功能。GreptimeDB v1.0 将其接入 Mito 扫描层。…

X AI KOLs Following 产品

摘要

GreptimeDB v1.0 集成了 DataFusion 的动态过滤下推功能,通过将运行时边界推送到扫描层来加速 TopK 查询,在 50 亿行的 trace 表上将查询时间从 29 秒降低到 0.21 秒。

DataFusion 在去年九月向上游合并了动态过滤下推功能。GreptimeDB v1.0 将其接入 Mito 扫描层。 当前 TopK 持有的边界会被下推为扫描层的运行时谓词——动态意味着它会随着 TopK 的收敛而收紧,并且每个行组都与最新的快照进行评估,而不是查询开始时的快照。 副作用:ORDER BY end_time DESC LIMIT 10(50亿行 traces 表)从 29 秒降至 0.21 秒。 https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter…
查看原文
查看缓存全文

缓存时间: 2026/05/19 12:47

DataFusion 在去年九月将动态过滤条件下推功能合并到了上游。GreptimeDB v1.0 将其接入到 Mito 扫描层。TopK 在运行过程中收紧的边界被作为运行时谓词下推到扫描层——动态的意思是,这个边界会随着 TopK 收敛而不断收紧,每个行组都会根据最新的快照来评估,而不是查询启动时的快照。效果:在 50 亿行的 traces 表上执行 ORDER BY end_time DESC LIMIT 10 从 29 秒降到了 0.21 秒。https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter… — # 从 29 秒到 0.21 秒:将 TopK 边界下推到 GreptimeDB 的扫描层 来源:https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter ORDER BY \.\.\. LIMIT 10 看起来应该很便宜,但通常并非如此。当表很大且排序键不是时间索引时,数据库通常需要读取大量行才能知道要返回哪十行。GreptimeDB v1.0 (https://www.greptime.com/blogs/2026-04-14-greptimedb-v1-ga-release) 通过将 TopK 算子正在构建的运行时边界尽早交还给扫描端来解决这个问题。一旦一个行组不可能对最终结果有贡献,扫描就不会读取它。在一个真实的 trace 数据集上,这个改变将 ORDER BY end_time DESC LIMIT 10 从大约半分钟降低到亚秒级。令人惊讶的是:同样的路径也取代了时间索引查询上旧的窗口排序 TopK 路径,因为更简单的 SortExec: TopK 加上动态过滤器方案测量起来更快。 这是分两步实现的。GreptimeDB #7545 (https://github.com/GreptimeTeam/greptimedb/pull/7545) 引入了核心的动态过滤器机制,将 DataFusion 上游的运行时 TopK 过滤 (https://datafusion.apache.org/blog/2025/09/10/dynamic-filters/) 接入到 GreptimeDB 自己的 Mito 扫描层。#7912 (https://github.com/GreptimeTeam/greptimedb/pull/7912) 随后将时间索引 TopK 的情况也归并到同一个路径上。本文的其余部分将介绍这个机制的工作原理以及它实际带来的好处。 — ## 为什么慢 (https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter#why-it-was-slow) sql SELECT start_time, end_time, run_type, status FROM langchain_traces ORDER BY end_time DESC LIMIT 10; 在旧路径上,扫描节点无法感知 TopK 的中间状态。尽管查询只返回 10 行,但扫描会一直读取直到查询结束。到所有操作完成时,扫描的大部分数据都是不必要的。这里有一个细节需要注意:start_time 是表的时间索引列,而 end_time 不是。GreptimeDB 之前针对时间索引 TopK 有一个优化叫做 窗口排序(内部使用 PartSortExec 进行本地排序)。它依赖于时间索引的数据分布,相当复杂。对于像 end_time 这样的非时间索引排序键,之前没有类似的裁剪能力。这就是 #7545 要填补的空白。#7912 则进一步推进了。基准测试表明,对于 TopK 来说,普通的 SortExec: TopK 加上动态过滤器WindowedSortExec 加上 PartSortExec 更快。因此,GreptimeDB 现在会在查询带有 LIMIT k 时禁用窗口排序重写,将时间索引和非时间索引的 TopK 都路由到相同的动态过滤器路径。窗口排序只保留给真正的全排序(没有 LIMIT)使用。 — ## 机制如何工作 (https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter#how-the-mechanism-works) 思路很直接:让 TopK 将其当前的边界作为运行时谓词反馈回扫描节点。 1. 在查询开始时,过滤器是敞开的,几乎允许所有数据通过。 2. 随着候选行的到达,当前 top-K 中最差的那一行成为新的边界——并且这个边界会不断收紧。 3. 该边界被封装成一个 DynamicFilterPhysicalExpr,并通过 DataFusion 的过滤器下推机制下推到扫描节点。 4. 扫描端将这个运行时条件与 SST 文件和行组的最小/最大统计信息结合起来,跳过那些不再可能符合条件的数据。 对于 ORDER BY end_time DESC LIMIT 10,胜利不在于结果集更小——始终只有十行需要返回。胜利在于 扫描范围在查询运行过程中不断缩小。下面的动画将两条路径并排展示。上半部分是旧路径:扫描层不知道 TopK 的阈值,因此所有八个行组都被读取。下半部分是新路径:阈值作为动态条件被下推,扫描器利用每个行组的最大值统计信息来跳过那些不可能进入前十的行组。 旧路径 vs 新路径:TopK 边界是否下推到扫描层 在执行计划层面,机制看起来像这样: ┌────────────────────────────────┐ │ TopK 持续更新其边界 │ └──────────────┬─────────────────┘ │ ▼ ┌────────────────────────────────┐ │ 动态过滤器收紧 │ └──────────────┬─────────────────┘ │ ▼ ┌────────────────────────────────┐ │ 扫描重新评估剩余的文件/行组 │ │ 使用行组统计信息 │ └──────────────┬─────────────────┘ │ ▼ ┌────────────────────────────────┐ │ 不符合条件的行组被完全跳过 │ └────────────────────────────────┘ — ## 边界的样子 (https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter#what-the-bound-looks-like) 一旦 TopK 构建了一个有意义的边界,GreptimeDB 会将其包装成一个对 NULL 敏感的谓词。对于我们的示例查询,它大致如下: text end_time IS NULL OR end_time > 1753660799999000000 IS NULL 分支的存在是因为 ORDER BY \.\.\. DESC NULLS FIRST 将 NULL 放在最前面,所以它们必须保持候选状态,不能被裁剪。 扫描端不需要读取每一行来获益。只要一个行组的统计信息表明其值范围不可能满足最新的条件,GreptimeDB 就可以在读取任何数据之前跳过整个行组。 — ## 确认它确实已接入 (https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter#confirming-it-s-actually-wired-up) EXPLAIN ANALYZE VERBOSE 展示了你需要的一切。TopK 节点显示其当前的运行时过滤器: text SortExec: TopK(fetch=10), expr=[end_time@1 DESC], preserve_partitioning=[true], filter=[end_time@1 IS NULL OR end_time@1 > 1753660799999000000] 扫描节点显示过滤器已被下推到它: text SeqScan: region=..., {"projection": [...], "dyn_filters": ["DynamicFilter [ end_time@1 IS NULL OR end_time@1 > 1753660799999000000 ]"], "files": [...]} 这两个信号都很重要。SortExec: TopK 携带 filter=[\.\.\.] 意味着 TopK 已经产生了运行时阈值;扫描节点上的 dyn_filters 意味着扫描层已接收到它并据此进行裁剪。当两者都出现时,动态过滤是端到端的。仓库中有对应的 sqlness 测试,位于 tests/cases/standalone/common/filter/topk\_dyn\_filter\.sql。 — ## 它带来的好处 (https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter#what-it-buys-you) 数据集是一个约 50-60 亿行的 langchain traces 表。最能说明问题的查询是本文开头的那个: sql SELECT * FROM langchain_traces ORDER BY end_time DESC LIMIT 10; ### 端到端延迟 (https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter#end-to-end-latency) | 查询 | 旧路径 | 动态过滤器路径 | 备注 | |——|––––|––––––––|——| | ORDER BY end_time DESC LIMIT 10 | ~28.9s | ~0.21s | end_time 不是时间索引;旧路径基本上全表扫描。有了动态过滤,扫描被更早地裁剪 | | ORDER BY start_time DESC LIMIT 10 | ~0.33s (窗口排序) | ~0.23–0.24s (动态过滤器) | start_time 是时间索引,原本由窗口排序处理。#7912 测量发现窗口排序在这里更慢,并将 TopK 移到了统一的动态过滤器路径上 | ### 算子级分解 (https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter#operator-level-breakdown) | 指标 | main | dyn_filter | 加速比 | |——|––––|—————|––––| | 总查询时间(用户时间) | 28.70 s | 0.20 s | ~143× | | 扫描节点总开销 | 28.81 s | 0.20 s | ~144× | | 排序执行计算时间 | 6.55 s | 0.009 s | ~720× | | 过滤前扫描行数 | 高(全表扫描) | 接近 0(已裁剪) | 显著 | end_time 的情况之所以慢,是因为旧路径扫描了整个表。有了动态过滤,TopK 边界提前裁剪了剩余的扫描——SortExec 自身的计算时间从 6.55 秒下降到 9 毫秒,仅仅是因为它不再需要比较那么多行。 start_time 的情况则不同。它以前走的是窗口排序(内部是 WindowedSortExecPartSortExec),依赖于分区范围和原生 SST 时间范围裁剪。在 #7545 时代,这个路径没有接入动态过滤器,但已经不错了(~0.33s)。#7912 进行了另一轮测量,并将其也移到了动态过滤器路径上,将 p50 进一步降低到约 0.23s。 明说一点:这个优化改变的是 扫描的数据量,而非算子复杂度。实际的加速效果取决于数据分布、行组统计信息的选择性以及 LIMIT 的大小。 — ## 关键代码路径 (https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter#key-code-paths) 代码库中有几个地方承载了这种端到端的能力: - RegionScanner trait (src/store\-api/src/region\_engine\.rs) 新增了一个方法: rust fn add_dyn_filter_to_predicate( &mut self, filter_exprs: Vec>, ) -> Vec; 返回值告诉 DataFusion 哪些过滤器实际被接受用于行组裁剪。未被接受的那些仍会在更高层级进行评估。 - RegionScanExec::handle\_child\_pushdown\_result (src/table/src/table/scan\.rs) 是 DataFusion 与 GreptimeDB 扫描层交汇的地方。它将来自父算子(TopK 或 Hash Join)的 parent_filters 交给扫描器,并向 DataFusion 的 FilterPushdownPropagation 报告支持情况。 - Predicate (src/table/src/predicate\.rs) 使用可热替换的结构存储动态过滤器: rust dyn_filters: Arc>>>, 通过 ArcSwap,扫描可以无锁地刷新边界。正在进行的读取最多只会看到一个过时的快照——它们不会阻塞,也不会破坏一致性。 - 所有三个 Mito 扫描路径——seq\_scan\.rsseries\_scan\.rsunordered\_scan\.rs——都将传入的过滤器转发到 ScanInput 内部的 PredicateGroup::add\_dyn\_filters,并填充 predicate\_allpredicate\_without\_region。 - 实际的裁剪发生在 FileRange::in\_dynamic\_filter\_range (src/mito2/src/sst/parquet/file\_range\.rs)。在读取每个行组之前,它会针对最新的 dyn_filters 结合 RowGroupPruningStats 运行 prune_with_stats,如果命中则跳过整个行组。 这个链条的关键特性是:每个行组都是根据最新的边界来评估的,而不是根据查询启动时获取的快照。 — ## 为什么时间索引 TopK 现在也使用动态过滤器 (https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter#why-time-index-topk-now-uses-dyn-filter-too) 在 #7545 中,刻意做出的选择是:不要为时间索引 TopK 启用动态过滤器下推。当时,PartSortExec 加动态过滤器的测量结果比单独使用任何一个路径都要慢,因此窗口排序保留了对这种情况的所有权。 #7912 重新审视了这个选择。发生了两件事。 **首先,PartSortExec::with\_new\_children 中的一个生命周期错误得到了修复。**旧代码先经过了 try_new,它构造了一个全新的 PartSortExec,丢弃了外部传入的动态过滤器句柄。实际上,动态过滤器从未看到 TopK 的当前边界。修复后,with_new_children 克隆现有的 PartSortExec,只交换输入,保留活动的动态过滤器引用: rust // 之前:try_new 重建了节点并给它一个没人更新的全新动态过滤器 // 之后:保留自身的动态过滤器,只更新输入 let mut new_exec = self.as_ref().clone(); new_exec.input = new_input.clone(); new_exec.properties = new_input.properties().clone(); **其次,在修复后重新运行了基准测试。**在 PartSortExec 加动态过滤器真正生效后,ORDER BY start_time DESC LIMIT k 的测量结果如下: | LIMIT | p50 | |—––|—–| | 1 | 0.586s | | 10 | 0.587s | | 100 | 0.569s | | 1000 | 0.584s | 一个 完全跳过窗口排序,回退到普通 SortExec: TopK 加动态过滤器 的变体则落在 0.23–0.24s 范围内。因此,即使有了正确性修复,窗口排序的 TopK 路径仍然比普通 TopK 加动态过滤器慢大约 2.5 倍。实际的代码改动只有一行加一个注释: rust // src/query/src/optimizer/windowed_sort.rs if /* ... 匹配窗口排序模式 ... */ && sort_exec.fetch().is_none() // 如果有 limit 则跳过,因为仅动态过滤器在这种情况下已经足够 { // 执行重写 } else { return Ok(Transformed::no(plan)); } 其行为变为: - SortExec 没有 fetch(即全排序,无 LIMIT)——仍然走窗口排序重写,保留分区范围裁剪; - SortExec 带有 fetch(即 TopK)——保持为 SortExec: TopK,让动态过滤器下推接手。 时间索引和非时间索引的排序键共享同一个路径。#7545 保留的 PartSortExec 例外现在已过时:在今天的执行计划中,TopK 根本不会产生 PartSortExec。 — ## 何时效果最佳 (https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter#when-this-helps-most) 最明显的收益发生在以下几个条件 同时满足 时: - k 很小,这样边界会迅速收紧; - 行组的最小/最大统计信息具有足够的选择性; - 数据分布使得 TopK 边界能尽早变得有意义。 排序键是否是时间索引过去是决定性因素,现在不再是了——TopK 无论哪种情况都使用相同的动态过滤器路径。唯一的区别是时间索引数据是天然有序的,因此边界收敛得稍快一些。 收益较小的情况: - 查询本身已经很快(例如,时间索引且窗口较小); - 过滤器在大部分运行时间内保持在接近 true 的状态; - 行组的最小/最大范围很宽,导致几乎无法裁剪; - LIMIT 足够大,边界永远不会变得紧凑。 简短版本:如果一个查询被扫描开销主导且数据形态配合,那么这个路径可以大幅改变成本曲线。 — ## 状态和下一步 (https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter#status-and-what-s-next) 由 #7545 和 #7912 实现的动态过滤当前覆盖: - 本地 TopKSortExec: TopK)运行时过滤器下推——已上线,现在对时间索引和非时间索引排序键统一。 - 本地 Hash Join(同一 Datanode 内的 HashJoin)动态过滤器——由 sqlness 测试案例覆盖,包括 hash\_join\_dyn\_filterhash\_join\_topk\_dyn\_filter。 - 分布式 Hash Join(Frontend 上的 join 将动态过滤器下推到 Datanode 上的扫描)——尚未支持。分布式版本需要一个“远程动态过滤器”机制:能够将在执行过程中产生的运行时条件通过 RPC 带回 Datanode 上的扫描算子。RFC #7931 和 #7979 中的基础设施工作仍在规划中。 — ## 故障排查清单 (https://greptime.com/blogs/2026-05-15-greptimedb-topk-dynamic-filter#troubleshooting-checklist) 如果你有一个缓慢的 ORDER BY \.\.\. LIMIT k 查询,请检查以下几点: 1. 计划中的 TopK 节点是否携带 filter=[\.\.\.] - 如果没有,查询可能根本就不是 TopK 形态——计划是全排序,或者 limit 发生在聚合之后。也可能你使用的 GreptimeDB 版本早于 #7912,此时时间索引列仍会被计划为 WindowedSortExec。 2. 扫描节点是否携带 dyn_filters: [\.\.\.] - 如果没有,动态过滤器没有成功下推。检查排序列、投影,以及中间是否有任何东西阻挡它(复杂表达式、UNION 等)。 3. 行组统计信息是否具有足够的选择性? - 如果行组的最小/最大范围很宽,裁剪效果可能不佳。使用 EXPLAIN ANALYZE 检查实际跳过的情况。 4. 能否在本地重现? - 使用 EXPLAIN VERBOSE 确认计划修改是否生效。如果动态过滤器没有出现,尝试简化查询(移除 JOIN、子查询)以隔离问题。

相似文章

@Greptime: GreptimeDB v1.1.2 已发布 — 这是一个值得升级的 v1.1 补丁。主要修复:定时 Flows 现在绑定 now()/current_timest…

X AI KOLs Following

GreptimeDB v1.1.2 是一个补丁版本,修复了定时 Flows 中 now() 的绑定问题,以确保 EVAL INTERVAL 窗口的确定性。此外还修复了以下问题:Kafka SASL 密码在调试输出中的屏蔽、GC 索引文件列表、parquet 元数据缓存大小、Prometheus 标签发现扫描以及 PromQL 时间二元聚合。建议用户升级。