From bf03399390876bd5233074b954ae34563639221d Mon Sep 17 00:00:00 2001 From: Kai Huang Date: Mon, 10 Aug 2026 11:13:27 -0700 Subject: [PATCH] Skip dedup-simplify rule on the analytics route PPL `dedup` followed by any pipe stage failed with a 500 on the analytics route: `IllegalStateException: Project rule encountered unmarked child [LogicalDedup]`. `dedup` is planned as a ROW_NUMBER() OVER (PARTITION BY keys) + Filter composite. PPLSimplifyDedupRule (run during CalciteToolsHelper.optimize) collapses that composite into a LogicalDedup so DedupPushdownRule can push it into a Lucene scan. The analytics engine has no such pushdown and cannot plan a LogicalDedup, so the node survives to its marking phase and fails. Add CalciteToolsHelper.optimizeForAnalytics, which runs the same rules minus PPLSimplifyDedupRule, and call it from the analytics route in RestUnifiedQueryAction. The analytics engine then receives the standard ROW_NUMBER window form directly. The Lucene/Calcite path is unchanged. Resolves #22671. Signed-off-by: Kai Huang --- .../sql/calcite/utils/CalciteToolsHelper.java | 18 ++++++++++ .../plugin/rest/RestUnifiedQueryAction.java | 2 +- .../sql/ppl/calcite/CalcitePPLDedupTest.java | 35 +++++++++++++++++++ 3 files changed, 54 insertions(+), 1 deletion(-) diff --git a/core/src/main/java/org/opensearch/sql/calcite/utils/CalciteToolsHelper.java b/core/src/main/java/org/opensearch/sql/calcite/utils/CalciteToolsHelper.java index 682a3ea17c1..c4a1ff1ac81 100644 --- a/core/src/main/java/org/opensearch/sql/calcite/utils/CalciteToolsHelper.java +++ b/core/src/main/java/org/opensearch/sql/calcite/utils/CalciteToolsHelper.java @@ -556,10 +556,28 @@ public RelNode visit(TableScan scan) { private static final HepProgram HEP_PROGRAM = new HepProgramBuilder().addRuleCollection(hepRuleList).build(); + // PPLSimplifyDedupRule collapses the ROW_NUMBER window form of dedup into a LogicalDedup so + // DedupPushdownRule can push it into a Lucene scan. The analytics engine has no such pushdown and + // cannot plan a LogicalDedup, so its optimization runs the standard window form directly. + private static final HepProgram ANALYTICS_HEP_PROGRAM = + new HepProgramBuilder() + .addRuleCollection( + hepRuleList.stream() + .filter(rule -> rule != PPLSimplifyDedupRule.DEDUP_SIMPLIFY_RULE) + .toList()) + .build(); + public static RelNode optimize(RelNode plan, CalcitePlanContext context) { Util.discard(context); HepPlanner planner = new HepPlanner(HEP_PROGRAM); planner.setRoot(plan); return planner.findBestExp(); } + + public static RelNode optimizeForAnalytics(RelNode plan, CalcitePlanContext context) { + Util.discard(context); + HepPlanner planner = new HepPlanner(ANALYTICS_HEP_PROGRAM); + planner.setRoot(plan); + return planner.findBestExp(); + } } diff --git a/plugin/src/main/java/org/opensearch/sql/plugin/rest/RestUnifiedQueryAction.java b/plugin/src/main/java/org/opensearch/sql/plugin/rest/RestUnifiedQueryAction.java index 609b4a71099..531d180f7cc 100644 --- a/plugin/src/main/java/org/opensearch/sql/plugin/rest/RestUnifiedQueryAction.java +++ b/plugin/src/main/java/org/opensearch/sql/plugin/rest/RestUnifiedQueryAction.java @@ -234,7 +234,7 @@ private void doExecute( plan = addFetchSizeLimit(plan, planContext, fetchSize); plan = addQuerySizeLimit(plan, planContext); plan = - org.opensearch.sql.calcite.utils.CalciteToolsHelper.optimize( + org.opensearch.sql.calcite.utils.CalciteToolsHelper.optimizeForAnalytics( plan, planContext); RelNode finalPlan = plan; Runnable executeTask = diff --git a/ppl/src/test/java/org/opensearch/sql/ppl/calcite/CalcitePPLDedupTest.java b/ppl/src/test/java/org/opensearch/sql/ppl/calcite/CalcitePPLDedupTest.java index 3c13297f8f0..5cd0fbe2568 100644 --- a/ppl/src/test/java/org/opensearch/sql/ppl/calcite/CalcitePPLDedupTest.java +++ b/ppl/src/test/java/org/opensearch/sql/ppl/calcite/CalcitePPLDedupTest.java @@ -14,6 +14,7 @@ import org.junit.Assert; import org.junit.Test; import org.opensearch.sql.calcite.plan.rule.PPLSimplifyDedupRule; +import org.opensearch.sql.calcite.utils.CalciteToolsHelper; public class CalcitePPLDedupTest extends CalcitePPLAbstractTest { @@ -487,4 +488,38 @@ private static RelNode runHepPlanner(RelNode root, HepProgram program) { planner.setRoot(root); return planner.findBestExp(); } + + /** + * {@code optimize} collapses dedup into a {@code LogicalDedup} (for Lucene pushdown), which the + * analytics engine cannot plan. {@code optimizeForAnalytics} omits that rule so dedup keeps the + * {@code ROW_NUMBER() OVER (PARTITION BY ...)} window form the analytics engine can plan + * (#22671). + */ + @Test + public void testOptimizeForAnalyticsKeepsRowNumberForm() { + RelNode raw = getRelNodeRaw("source=EMP | where SAL > 1000 | dedup DEPTNO"); + + Assert.assertTrue( + "optimize() should collapse dedup into a LogicalDedup:\n", + CalciteToolsHelper.optimize(raw, null).explain().contains("LogicalDedup")); + + String analytics = CalciteToolsHelper.optimizeForAnalytics(raw, null).explain(); + Assert.assertFalse( + "optimizeForAnalytics should not produce a LogicalDedup:\n" + analytics, + analytics.contains("LogicalDedup")); + Assert.assertTrue( + "optimizeForAnalytics should keep the ROW_NUMBER window form:\n" + analytics, + analytics.contains("ROW_NUMBER")); + Assert.assertTrue( + "where predicate should be preserved:\n" + analytics, analytics.contains(">($5, 1000)")); + } + + /** A dedup-free plan is identical under both optimize variants. */ + @Test + public void testOptimizeForAnalyticsMatchesOptimizeWithoutDedup() { + RelNode raw = getRelNodeRaw("source=EMP | where SAL > 1000 | fields DEPTNO"); + Assert.assertEquals( + CalciteToolsHelper.optimize(raw, null).explain(), + CalciteToolsHelper.optimizeForAnalytics(raw, null).explain()); + } }