Fix ordering propagation around table and window functions - #18749
Conversation
Add complete outer ORDER BY clauses to ordered result assertions, including rows with equal values and device IDs. Resolve input ties with time only for position-sensitive windows, preserving peer-based window semantics and all existing expected results. Verified both window function IT classes on Java 17 with TableClusterIT: 52 tests passed against the baseline engine without the separate TVF fix.
Propagate proven ordering through known order-preserving window TVFs and partition keys, preserving sort elimination where safe. Retain sorts for generated columns and custom row-semantic functions, and require a proven streaming prefix before selecting StreamSort. Cover distributed scan ordering, generated-column ties, set and custom row semantics, nested functions and blocked ordering propagation with planner and integration regressions.
Align the filter-pushdown expected rows with ORDER BY device, time. Use time to break equal-value row_number ties in both filter and limit pushdown tests so the selected rows and their numbers are deterministic. Verified all 13 IoTDBWindowFunction3IT cases under both TableSimpleIT and TableClusterIT on Java 17 (26 passing cases total).
Wait for complete groups before producing stream-sort output, propagate TopK input readiness, and copy retained rows when compacting TopK buffers. Add regressions for small output blocks, merge-input progress, and retained row values.
Derive streaming aggregation from a valid partition prefix and propagate actual ordering through window, ranking, pattern, and table-function plans. Keep native device/time ordering and streaming sorts without dropping regional inputs or treating per-partition order as global order. Preserve sort restrictions across FILL, account for nullable filled keys and DESC null placement, and add planner and integration regressions for correctness and sort elimination.
| if (!(node instanceof ValueFillNode)) { | ||
| context.clearExpectedOrderingScheme(); | ||
| List<PlanNode> result = dealWithPlainSingleChildNode(node, context); | ||
| nodeOrderingMap.remove(node.getPlanNodeId()); |
There was a problem hiding this comment.
Perf regression: ORDER BY <group keys>, time over FILL METHOD PREVIOUS | NEXT | LINEAR ... FILL_GROUP now gets an extra blocking sort above the fill
The non-constant branch of visitFill now drops the fill's ordering unconditionally (this line). With FILL_GROUP, QueryPlanner.fillGroup already sorts the fill input by (group keys, TIME_COLUMN), so on master an outer ORDER BY <group keys>, time is eliminated. With this PR it becomes a full SortNode in the root fragment, which buffers the whole result (and may spill) before emitting anything, for the common "fill per device, return per device" query:
SELECT time, tag1, tag2, tag3, s1 FROM table1
FILL METHOD PREVIOUS TIME_COLUMN 1 FILL_GROUP 2,3,4
ORDER BY tag1, tag2, tag3, timemaster (and with the change below) this PR (20d7b00)
Output Output
PreviousFill Sort [tag1, tag2, tag3, time] <-- new
MergeSort [tag1, tag2, tag3, time] PreviousFill
DeviceTableScan x 3 regions MergeSort [tag1, tag2, tag3, time]
DeviceTableScan x 3 regions
Same pattern for these shapes (distributed plans from PlanTester, 3 regions):
| Shape | master | this PR | with the change below |
|---|---|---|---|
LINEAR TIME_COLUMN 1 FILL_GROUP 2,3,4 or NEXT TIME_COLUMN 1 FILL_GROUP 2,3,4, ORDER BY tag1, tag2, tag3, time |
no sort above the fill | + full Sort | as master |
PREVIOUS TIME_BOUND 1h TIME_COLUMN 1 FILL_GROUP 2,3,4, same ORDER BY |
no sort above the fill | + full Sort | as master |
PREVIOUS TIME_COLUMN 1 FILL_GROUP 2,3,4 ORDER BY tag1, tag2, tag3 |
no sort above the fill | + full Sort | as master |
SELECT time, tag1, s1 ... FILL_GROUP 2 ORDER BY tag1, time |
no sort above the fill | + full Sort | as master |
PREVIOUS ... FILL_GROUP 2,3,4 over HOP(DATA => table1, ...) |
no sort above the fill | + full Sort | as master |
PREVIOUS ... FILL_GROUP 2,3,4 with lag(s1) OVER (PARTITION BY tag1, tag2, tag3 ORDER BY time) in the select list |
StreamSort below the fill | + full Sort above the fill | no sort at all |
date_bin_gapfill(...) or date_bin(...) aggregation + FILL ... FILL_GROUP 2,3,4 ORDER BY tag1, tag2, tag3, t |
1 Sort (below the fill) | 2 Sorts | still 2 Sorts (see the end) |
Dropping the ordering is right for sort keys the fill can rewrite, and it fixes a real master bug: master eliminates the outer sort in SELECT * FROM (SELECT * FROM table1 ORDER BY s1, s2 LIMIT 10) FILL METHOD PREVIOUS ORDER BY s1, s2 (also with NEXT / LINEAR), although PREVIOUS turns (5, 1), (5, 9), (NULL, 2) into (5, 1), (5, 9), (5, 2). But part of the child ordering provably survives a PREVIOUS / NEXT / LINEAR fill:
- the fill only writes NULL positions and emits rows in input order;
- it resets at every FILL_GROUP boundary, and all rows of a group share their group-key values, so a group key is never rewritten (a NULL key has no non-NULL value in its group to copy);
- non-NULL columns (native time, tags / attributes that are non-null for every device) are never rewritten;
getNonNullSymbolsalready proves this for the constant fill.
So the fill can keep the prefix of its child's ordering that consists of FILL_GROUP keys and provably non-null columns. As in the constant-fill branch, getNonNullSymbols (which may read the coordinator's spilled device entries) only runs when an ancestor requests an ordering; the group keys alone need no metadata.
--- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/distribute/TableDistributedPlanGenerator.java
+++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/planner/distribute/TableDistributedPlanGenerator.java
@@ -52,11 +52,14 @@
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.GroupNode;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.JoinNode;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.LimitNode;
+import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.LinearFillNode;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.MarkDistinctNode;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.MergeSortNode;
+import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.NextFillNode;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.OffsetNode;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.OutputNode;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.PatternRecognitionNode;
+import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.PreviousFillNode;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.ProjectNode;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.RowNumberNode;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.SemiJoinNode;
@@ -323,9 +326,24 @@
@Override
public List<PlanNode> visitFill(FillNode node, PlanContext context) {
if (!(node instanceof ValueFillNode)) {
+ boolean hasExpectedOutputOrdering = context.hasSortProperty;
context.clearExpectedOrderingScheme();
+ // PREVIOUS/NEXT/LINEAR fill only replaces NULLs, keeps the row order and resets at every
+ // FILL_GROUP boundary. Group keys are constant within a group and non-null columns are never
+ // rewritten, so a child ordering prefix made of such columns still holds after the fill.
+ Set<Symbol> unchanged = new HashSet<>(getFillGroupingKeys(node));
+ if (hasExpectedOutputOrdering) {
+ unchanged.addAll(getNonNullSymbols(node.getChild()));
+ }
List<PlanNode> result = dealWithPlainSingleChildNode(node, context);
- nodeOrderingMap.remove(node.getPlanNodeId());
+ OrderingScheme ordering =
+ TableFunctionOrdering.retainOrderingPrefix(
+ nodeOrderingMap.get(node.getPlanNodeId()), unchanged);
+ if (ordering == null) {
+ nodeOrderingMap.remove(node.getPlanNodeId());
+ } else {
+ nodeOrderingMap.put(node.getPlanNodeId(), ordering);
+ }
return result;
}
// Inspect metadata before distribution consumes a spilled device data set. No data rows need
@@ -354,6 +372,18 @@
return Collections.singletonList(node);
}
+ private static List<Symbol> getFillGroupingKeys(FillNode node) {
+ Optional<List<Symbol>> groupingKeys = Optional.empty();
+ if (node instanceof PreviousFillNode) {
+ groupingKeys = ((PreviousFillNode) node).getGroupingKeys();
+ } else if (node instanceof NextFillNode) {
+ groupingKeys = ((NextFillNode) node).getGroupingKeys();
+ } else if (node instanceof LinearFillNode) {
+ groupingKeys = ((LinearFillNode) node).getGroupingKeys();
+ }
+ return groupingKeys.orElse(Collections.emptyList());
+ }
+
private Set<Symbol> getNonNullSymbols(PlanNode node) {
Set<Symbol> cached = nonNullSymbols.get(node.getPlanNodeId());
if (cached != null) {Regression tests (OrderingPropertyPropagationTest)
--- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/OrderingPropertyPropagationTest.java
+++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/relational/planner/optimizations/OrderingPropertyPropagationTest.java
@@ -23,6 +23,7 @@
import org.apache.iotdb.commons.queryengine.plan.relational.planner.Symbol;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.AggregationNode;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.CollectNode;
+import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.FillNode;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.PatternRecognitionNode;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.RowNumberNode;
import org.apache.iotdb.commons.queryengine.plan.relational.planner.node.SortNode;
@@ -296,6 +297,48 @@
}
@Test
+ public void testFillGroupKeepsGroupKeyAndTimeOrdering() {
+ String window = ", lag(s1) OVER (PARTITION BY " + TAGS + " ORDER BY time) AS lg";
+ for (String extraColumn : List.of("", window)) {
+ for (String method :
+ List.of(
+ "PREVIOUS TIME_COLUMN 1",
+ "PREVIOUS TIME_BOUND 1h TIME_COLUMN 1",
+ "NEXT TIME_COLUMN 1",
+ "LINEAR TIME_COLUMN 1")) {
+ for (String[] group : new String[][] {{"2,3,4", TAGS}, {"2", "tag1"}}) {
+ String sql =
+ "SELECT time, "
+ + TAGS
+ + ", s1"
+ + extraColumn
+ + " FROM table1 FILL METHOD "
+ + method
+ + " FILL_GROUP "
+ + group[0]
+ + " ORDER BY "
+ + group[1]
+ + ", time";
+ PlanTester tester = plan(sql);
+ assertFalse(sql, findNodes(tester, FillNode.class).isEmpty());
+ // FILL_GROUP already orders the input by its keys and time, and the fill changes neither.
+ assertTrue(sql, sortsAbove(tester, FillNode.class).isEmpty());
+ }
+ }
+ }
+ }
+
+ @Test
+ public void testFillKeepsSortOnFilledNullableKeys() {
+ // PREVIOUS can turn (5, 1), (5, 9), (NULL, 2) into (5, 1), (5, 9), (5, 2), so an ordering on
+ // nullable columns that are not FILL_GROUP keys must not survive the fill.
+ String sql =
+ "SELECT * FROM (SELECT * FROM table1 ORDER BY s1, s2 LIMIT 10)"
+ + " FILL METHOD PREVIOUS ORDER BY s1, s2";
+ assertFalse(sql, sortsAbove(plan(sql), FillNode.class).isEmpty());
+ }
+
+ @Test
public void testUnorderedQueriesKeepUnorderedCollection() {
for (String sql :
List.of(Checked locally on top of 20d7b00:
testFillGroupKeepsGroupKeyAndTimeOrderingfails on 20d7b00 (Sort above the fill) and passes with the change.testFillKeepsSortOnFilledNullableKeyspasses on both, so the master bug above stays fixed.- Of 149 distributed-plan probes (TVF / window / TopK / fill / aggregation / LIMIT / JOIN / UNION / single-device shapes), only the 8 FILL_GROUP shapes above that the change fixes differ from 20d7b00; the other 141 plans are identical.
- All unit tests under
org/apache/iotdb/db/queryengine/plan/relationalpass (493 run) exceptCteMaterializerTest/CteSubqueryTest, which error locally on JDK 21 with or without the change.
Not covered: gapfill + FILL_GROUP and date_bin aggregation + FILL_GROUP still sort twice. There the bucket column comes out of an aggregation, which getNonNullSymbols does not model, and AbstractGapFillOperator explicitly passes NULL buckets through, so I did not assume it is non-NULL. Covering these needs a separate argument about how NULLs in the key right after the group keys get filled; I'd leave that for a follow-up.
Stop cached NEXT and LINEAR lookahead at group boundaries. Reset linear-fill state before preparing a new group's lookahead so retries retain valid next values. Cover cross-block gaps, null helpers, and prepared values with operator regressions.
Retain FILL_GROUP and proven non-null key prefixes. Preserve a chronological NULLS LAST helper prefix without assuming secondary peer order, avoiding redundant sorts for grouped fill and time-bucket aggregation. Add planner and integration regressions for PREVIOUS, NEXT, LINEAR, time bounds, nullable keys, windows, TVFs, and gapfill. Use the existing GapFill IT sort-buffer setting for parallel cluster validation.
Description
Queries that order TVF or window output can lose required sorting, fail when aggregating a non-prefix subset of TVF partition keys, or add blocking sorts to already ordered device/time streams. Derive output order from each operator's actual guarantees and retain only proven prefixes. Preserve native scan/merge order for TUMBLE, HOP, CUMULATE, SESSION, VARIATION, and CAPACITY; unknown TVFs retain only a proven partition prefix.
Keep safe sort elimination and streaming sort for common device/time, window, LIMIT, and TopK queries. Do not sort all TopK input rows by ranking fields. Preserve all regional inputs when Group has already been eliminated, and distinguish per-partition rank order from global order. Constant FILL retains order only for unaffected keys. PREVIOUS/NEXT/LINEAR preserve prefixes of FILL_GROUP keys and proven non-null columns; an ascending NULLS LAST time key immediately after all group keys is also preserved. A filled time key may merge peers, so secondary ordering is retained only when separately proven. This also removes redundant sorting after grouped date_bin/date_bin_gapfill aggregation. Sort-elimination restrictions cannot be cleared by another node or sibling.
The change also protects non-TVF queries whose requested streaming prefix is not supplied by a window's actual partitioning, DIFF, or TopK output, and corrects DESC NULL placement in native device ordering. Integration-test tie ordering is explicit, and timestamp/window-index assertions use typed comparisons.
Runtime verification with small TsBlocks also exposed three execution issues: stream sort flushed incomplete groups when a previous group left buffered output; TopK did not propagate merge-input readiness; and TopK compaction changed row positions without copying the underlying block. Fix these without replacing streaming sort with a blocking sort, and add operator regressions for progress and exact retained values.
Grouped NEXT/LINEAR lookahead stops at cached group boundaries. LINEAR resets state before preparing values for a new group, preserving valid lookahead across retries and preventing values from leaking between groups.
Validation
This PR has:
Key changed classes and packages
TableDistributedPlanGenerator,SortElimination,TableFunctionOrdering,TransformAggregationToStreamable,TransformSortToStreamSort,TableBuiltinTableFunction, andorg.apache.iotdb.calc.execution.operator.