Hi weihao,
nice to see such a great performance improvement for table model TopK query which is a high-frequency query in POC cases, and I've approved that pr, but there are still some CI problems, you need to fix them and after that I'll merged that. And since our 2.0.11 features have already frozen, I think we'll include your commit in the next version which will be 2.0.12. Best regards, --------------------- Yuan Tian On Wed, Jul 15, 2026 at 7:40 PM Weihao Li <[email protected]> wrote: > Hi all, > > > > This PR introduces TopK Runtime Filter for table-model ORDER BY time + > LIMIT queries. During execution, TopK bounds are used to filter unrelated > devices/data early, so we no longer need a full scan just to return limit > rows. On the logical side we mark eligible TopK/Scan nodes; in distributed > plans the source id is propagated to region-level TopKs. > > > Benchmark setup: 100,000 devices × 20 measurements, 5s sampling interval, > 100,000 rows per device. Theoretical speedup is about 100,000 / limit > (after optimization we only need to read up to limit devices instead of > everything). > > > Results: You can find detail results in PR descriptor, here is a brief. > | > LIMIT > | > Before > | > After > | > Actual Speedup > | > | > 10 / 100 / 1000 > | > 17–25s > | > 0.41–0.74s > | > 34 – 44× > | > | > 10000 / 50000 / 100000 > | > timeout ( more than 60s) > | > 6.7–51s > | > more than 1 - 9x > | > > > Small LIMITs drop from tens of seconds to sub-second latency; large LIMITs > go from timeout to finishing. Actual speedup is below the theoretical upper > bound (which assumes only limit devices are read), but the gain is still > clear. > > > Yours, > Weihao Li > > > > > > > > > > > > At 2026-07-10 18:23:30, "Weihao Li" <[email protected]> wrote: > >Hi all, > > > > > >I'd like to share a query optimization we are working on for the table > model: TopK Runtime Filter for full-table ORDER BY time LIMIT k queries. > >Typical queries: > >SELECT * FROM t1 ORDER BY time DESC LIMIT 10; > >SELECT * FROM t1 ORDER BY time ASC LIMIT 10; > > > >1. Current status > >Optimization path (before this change) > >Logical plan: MergeLimitWithSort merges Limit + Sort into TopKNode > >LIMIT pushdown: PushLimitOffsetIntoTableScan pushes LIMIT into > DeviceTableScanNode; with multiple devices, pushLimitToEachDevice = true > >Distributed plan: TableDistributedPlanGenerator splits Scan by Region; > MergeLimitWithMergeSort merges Limit + MergeSort into TopK; > AddExchangeNodes inserts Exchange > >Root TopK merge: each DataNode’s local TopK results are exchanged and > merged by the root TopK > > > > > >Problem > >When there are many devices (Tag dimension), each device returns up to k > rows, so Scan may read about k × device_count rows before the global TopK > keeps only k. When the device count is large, Scan I/O and CPU cost grow a > lot, even though the final result only needs k rows. > > > > > >2. Optimization idea > >We adopt a two-phase design, this helps to reduce read load a lot : > > > >Optimize- Add TopKRuntimeFilterOptimizer to mark if TopK/TableScan need > create or set RuntimeFilter. > >Execution-When generate operators, create a shared RuntimeFilter; TopK > updates the threshold; Scan prunes read path using RuntimeFilter. > > > >3. Details > >3.1 Plan-time pipeline > >In TableDistributedPlanner.generateDistributedPlanWithOptimize(): > >1. TableDistributedPlanGenerator.genResult() // split by Region > >2. DistributedOptimizeFactory optimizers // e.g. MergeLimitWithMergeSort > >3. AddExchangeNodes() // insert Exchange > >4. if enable_topk_runtime_filter == true: > > TopKRuntimeFilterOptimizer.optimize() // MUST run after Exchange > insertion > >5. SubPlanGenerator.splitToSubPlan() // split into Fragments > > > >Trigger conditions (all required): > >enable_topk_runtime_filter = true > >Single-column ORDER BY time > >TopK subtree contains a local DeviceTableScanNode > > > > > >Typical distributed plan: > >Fragment-0 (Root): > > TopK ← subtree has only Exchange → NOT marked > > ├── Exchange ── Fragment-1 > > └── Exchange ── Fragment-2 > > > > > >Fragment-1: > > TopK (TOPN OPT) ← local TopK + Scan → marked > > └── Limit > > └── DeviceTableScan (TOPN OPT: {topKId}) > > > >3.2 Runtime pipeline > > > >1. TableOperatorGenerator -> vistTopK -> produce RuntimeFilter and record > it into context, visitTableScan -> read RuntimeFilter from context and set > it into TableScanOperator > >2. TableTopKOperator -> when the matained heap of size k was full, call > updateThreshold(peek.time) of RuntimeFilter > >3. TableScanOperator -> use RuntimeFilter to prune read load > >** mayQualifyRange(start, end) — skip resources / files / chunks / pages > >** mayQualify(time) — row-level pruning > > > > > >Yours, > >Weihao Li > > >
