From 19732b34ded1dbee085d3e3f65ccd65c98166501 Mon Sep 17 00:00:00 2001 From: Heng Qian Date: Wed, 26 Feb 2025 10:54:23 +0800 Subject: [PATCH 1/5] Transform to calcite plan before executing Signed-off-by: Heng Qian --- .../opensearch/sql/executor/QueryService.java | 29 ++++++++++++++++++- 1 file changed, 28 insertions(+), 1 deletion(-) diff --git a/core/src/main/java/org/opensearch/sql/executor/QueryService.java b/core/src/main/java/org/opensearch/sql/executor/QueryService.java index 7a04ad7deff..73d6a9aa5d0 100644 --- a/core/src/main/java/org/opensearch/sql/executor/QueryService.java +++ b/core/src/main/java/org/opensearch/sql/executor/QueryService.java @@ -17,7 +17,10 @@ import lombok.RequiredArgsConstructor; import org.apache.calcite.jdbc.CalciteSchema; import org.apache.calcite.plan.RelTraitDef; +import org.apache.calcite.rel.RelCollation; +import org.apache.calcite.rel.RelCollations; import org.apache.calcite.rel.RelNode; +import org.apache.calcite.rel.logical.LogicalSort; import org.apache.calcite.rel.metadata.DefaultRelMetadataProvider; import org.apache.calcite.schema.SchemaPlus; import org.apache.calcite.sql.parser.SqlParser; @@ -26,6 +29,7 @@ import org.apache.calcite.tools.Programs; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; +import org.jetbrains.annotations.NotNull; import org.opensearch.sql.analysis.AnalysisContext; import org.opensearch.sql.analysis.Analyzer; import org.opensearch.sql.ast.tree.UnresolvedPlan; @@ -132,7 +136,30 @@ public void executePlanByCalcite( RelNode plan, CalcitePlanContext context, ResponseListener listener) { - executionEngine.execute(optimize(plan), context, listener); + executionEngine.execute(convertToCalcitePlan(optimize(plan)), context, listener); + } + + /** + * Convert OpenSearch Plan to Calcite Plan. + * Although both plans consist of Calcite RelNodes, + * there are some differences in the topological structures or semantics between them. + * + * @param osPlan Logical Plan derived from OpenSearch PPL + */ + private static @NotNull RelNode convertToCalcitePlan(RelNode osPlan) { + RelNode calcitePlan = osPlan; + + // Calcite only ensures collation of the final result produced from the root sort operator. + // While we expect that the collation can be preserved through the pipes over PPL, we need to + // explicitly add a sort operator on top of the original plan + // to ensure the correct collation of the final result. + // See logic in ${@link } + // For the redundant sort, we rely on Calcite optimizer to eliminate + RelCollation collation = osPlan.getTraitSet().getCollation(); + if (collation != RelCollations.EMPTY) { + calcitePlan = LogicalSort.create(osPlan, collation, null, null); + } + return calcitePlan; } /** From 378d1657c9322b51b5d4d757764d5e7b4d4e11ed Mon Sep 17 00:00:00 2001 From: Heng Qian Date: Wed, 26 Feb 2025 14:43:18 +0800 Subject: [PATCH 2/5] Fix bug for single column row Signed-off-by: Heng Qian --- .../opensearch/sql/executor/QueryService.java | 22 ++++++++++--------- .../scan/CalciteOpenSearchIndexScan.java | 10 +++++++-- .../scan/OpenSearchIndexEnumerator.java | 9 ++++++-- 3 files changed, 27 insertions(+), 14 deletions(-) diff --git a/core/src/main/java/org/opensearch/sql/executor/QueryService.java b/core/src/main/java/org/opensearch/sql/executor/QueryService.java index 73d6a9aa5d0..ca67bd3efce 100644 --- a/core/src/main/java/org/opensearch/sql/executor/QueryService.java +++ b/core/src/main/java/org/opensearch/sql/executor/QueryService.java @@ -20,6 +20,7 @@ import org.apache.calcite.rel.RelCollation; import org.apache.calcite.rel.RelCollations; import org.apache.calcite.rel.RelNode; +import org.apache.calcite.rel.core.Sort; import org.apache.calcite.rel.logical.LogicalSort; import org.apache.calcite.rel.metadata.DefaultRelMetadataProvider; import org.apache.calcite.schema.SchemaPlus; @@ -140,25 +141,26 @@ public void executePlanByCalcite( } /** - * Convert OpenSearch Plan to Calcite Plan. - * Although both plans consist of Calcite RelNodes, - * there are some differences in the topological structures or semantics between them. + * Convert OpenSearch Plan to Calcite Plan. Although both plans consist of Calcite RelNodes, there + * are some differences in the topological structures or semantics between them. * * @param osPlan Logical Plan derived from OpenSearch PPL */ private static @NotNull RelNode convertToCalcitePlan(RelNode osPlan) { RelNode calcitePlan = osPlan; - // Calcite only ensures collation of the final result produced from the root sort operator. - // While we expect that the collation can be preserved through the pipes over PPL, we need to - // explicitly add a sort operator on top of the original plan - // to ensure the correct collation of the final result. - // See logic in ${@link } - // For the redundant sort, we rely on Calcite optimizer to eliminate + /* Calcite only ensures collation of the final result produced from the root sort operator. + * While we expect that the collation can be preserved through the pipes over PPL, we need to + * explicitly add a sort operator on top of the original plan + * to ensure the correct collation of the final result. + * See logic in ${@link CalcitePrepareImpl} + * For the redundant sort, we rely on Calcite optimizer to eliminate + */ RelCollation collation = osPlan.getTraitSet().getCollation(); - if (collation != RelCollations.EMPTY) { + if (!(osPlan instanceof Sort) && collation != RelCollations.EMPTY) { calcitePlan = LogicalSort.create(osPlan, collation, null, null); } + return calcitePlan; } diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java index 20fb52c1125..3a105b64f89 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java @@ -97,9 +97,15 @@ public RelDataType deriveRowType() { @Override public Result implement(EnumerableRelImplementor implementor, Prefer pref) { - // Avoid optimizing the java row type since the scan will always return an array. + /* In Calcite enumerable operators, row of single column will be optimized to a scalar value. + * See {@link PhysTypeImpl}. + * Since we need to combine this operator with their original ones, + * let's follow this convention to apply the optimization here and ensure `scan` method + * returns the correct data format for single column rows. + * See {@link OpenSearchIndexEnumerator} + */ PhysType physType = - PhysTypeImpl.of(implementor.getTypeFactory(), getRowType(), pref.preferArray(), false); + PhysTypeImpl.of(implementor.getTypeFactory(), getRowType(), pref.preferArray()); Expression scanOperator = implementor.stash(this, CalciteOpenSearchIndexScan.class); return implementor.result(physType, Blocks.toBlock(Expressions.call(scanOperator, "scan"))); diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/OpenSearchIndexEnumerator.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/OpenSearchIndexEnumerator.java index 518e67d49fa..6e778422db2 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/OpenSearchIndexEnumerator.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/OpenSearchIndexEnumerator.java @@ -61,8 +61,13 @@ private void fetchNextBatch() { @Override public Object current() { - Object[] p = fields.stream().map(k -> current.tupleValue().get(k).valueForCalcite()).toArray(); - return p; + /* In Calcite enumerable operators, row of single column will be optimized to a scalar value. + * See {@link PhysTypeImpl} + */ + if (fields.size() == 1) { + return current.tupleValue().get(fields.getFirst()).valueForCalcite(); + } + return fields.stream().map(k -> current.tupleValue().get(k).valueForCalcite()).toArray(); } @Override From 4698600ce1c45928a091abc1aed5a2a57ca59b4b Mon Sep 17 00:00:00 2001 From: Heng Qian Date: Wed, 26 Feb 2025 15:11:19 +0800 Subject: [PATCH 3/5] Add settings for calcite pushdown Signed-off-by: Heng Qian --- .../java/org/opensearch/sql/common/setting/Settings.java | 1 + .../sql/calcite/standalone/CalcitePPLIntegTestCase.java | 1 + .../sql/opensearch/setting/OpenSearchSettings.java | 8 ++++++++ .../sql/opensearch/storage/OpenSearchIndex.java | 2 +- .../storage/scan/CalciteOpenSearchIndexScan.java | 7 +++++-- 5 files changed, 16 insertions(+), 3 deletions(-) diff --git a/common/src/main/java/org/opensearch/sql/common/setting/Settings.java b/common/src/main/java/org/opensearch/sql/common/setting/Settings.java index 63ee60b7683..6225d8fe0b7 100644 --- a/common/src/main/java/org/opensearch/sql/common/setting/Settings.java +++ b/common/src/main/java/org/opensearch/sql/common/setting/Settings.java @@ -31,6 +31,7 @@ public enum Key { /** Enable Calcite as execution engine */ CALCITE_ENGINE_ENABLED("plugins.calcite.enabled"), CALCITE_FALLBACK_ALLOWED("plugins.calcite.fallback.allowed"), + CALCITE_PUSHDOWN_ENABLED("plugins.calcite.pushdown.enabled"), /** Query Settings. */ FIELD_TYPE_TOLERANCE("plugins.query.field_type_tolerance"), diff --git a/integ-test/src/test/java/org/opensearch/sql/calcite/standalone/CalcitePPLIntegTestCase.java b/integ-test/src/test/java/org/opensearch/sql/calcite/standalone/CalcitePPLIntegTestCase.java index 9a366c2a2a0..4384bee0f48 100644 --- a/integ-test/src/test/java/org/opensearch/sql/calcite/standalone/CalcitePPLIntegTestCase.java +++ b/integ-test/src/test/java/org/opensearch/sql/calcite/standalone/CalcitePPLIntegTestCase.java @@ -104,6 +104,7 @@ private Settings defaultSettings() { .put(Key.FIELD_TYPE_TOLERANCE, true) .put(Key.CALCITE_ENGINE_ENABLED, true) .put(Key.CALCITE_FALLBACK_ALLOWED, false) + .put(Key.CALCITE_PUSHDOWN_ENABLED, false) .build(); @Override diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/setting/OpenSearchSettings.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/setting/OpenSearchSettings.java index 53bf6536c96..f1f12a21553 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/setting/OpenSearchSettings.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/setting/OpenSearchSettings.java @@ -99,6 +99,13 @@ public class OpenSearchSettings extends Settings { Setting.Property.NodeScope, Setting.Property.Dynamic); + public static final Setting CALCITE_PUSHDOWN_ENABLED_SETTING = + Setting.boolSetting( + Key.CALCITE_PUSHDOWN_ENABLED.getKeyValue(), + true, + Setting.Property.NodeScope, + Setting.Property.Dynamic); + public static final Setting QUERY_MEMORY_LIMIT_SETTING = new Setting<>( Key.QUERY_MEMORY_LIMIT.getKeyValue(), @@ -478,6 +485,7 @@ public static List> pluginSettings() { .add(PPL_ENABLED_SETTING) .add(CALCITE_ENGINE_ENABLED_SETTING) .add(CALCITE_FALLBACK_ALLOWED_SETTING) + .add(CALCITE_PUSHDOWN_ENABLED_SETTING) .add(QUERY_MEMORY_LIMIT_SETTING) .add(QUERY_SIZE_LIMIT_SETTING) .add(METRICS_ROLLING_WINDOW_SETTING) diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/OpenSearchIndex.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/OpenSearchIndex.java index 99df0465bdc..c77bb3c94d7 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/OpenSearchIndex.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/OpenSearchIndex.java @@ -68,7 +68,7 @@ public class OpenSearchIndex extends OpenSearchTable { /** OpenSearch client connection. */ @Getter private final OpenSearchClient client; - private final Settings settings; + @Getter private final Settings settings; /** {@link OpenSearchRequest.IndexName}. */ private final OpenSearchRequest.IndexName indexName; diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java index 3a105b64f89..54a350c7428 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java @@ -32,6 +32,7 @@ import org.checkerframework.checker.nullness.qual.Nullable; import org.opensearch.index.query.QueryBuilder; import org.opensearch.sql.calcite.plan.OpenSearchTableScan; +import org.opensearch.sql.common.setting.Settings; import org.opensearch.sql.opensearch.planner.physical.OpenSearchIndexRules; import org.opensearch.sql.opensearch.request.OpenSearchRequestBuilder; import org.opensearch.sql.opensearch.request.PredicateAnalyzer; @@ -85,8 +86,10 @@ public RelNode copy(RelTraitSet traitSet, List inputs) { @Override public void register(RelOptPlanner planner) { super.register(planner); - for (RelOptRule rule : OpenSearchIndexRules.OPEN_SEARCH_INDEX_SCAN_RULES) { - planner.addRule(rule); + if (osIndex.getSettings().getSettingValue(Settings.Key.CALCITE_PUSHDOWN_ENABLED)) { + for (RelOptRule rule : OpenSearchIndexRules.OPEN_SEARCH_INDEX_SCAN_RULES) { + planner.addRule(rule); + } } } From efc8cd8282b205c68262711045372c8d5aae9aa3 Mon Sep 17 00:00:00 2001 From: Heng Qian Date: Wed, 26 Feb 2025 16:21:21 +0800 Subject: [PATCH 4/5] Lazily construct OpenSearchRequestBuilder and do push down Signed-off-by: Heng Qian --- .../OpenSearchFilterIndexScanRule.java | 3 +- .../scan/CalciteOpenSearchIndexScan.java | 56 +++++++++++++------ 2 files changed, 42 insertions(+), 17 deletions(-) diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/planner/physical/OpenSearchFilterIndexScanRule.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/planner/physical/OpenSearchFilterIndexScanRule.java index f8beb339fe1..2b1e8e53124 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/planner/physical/OpenSearchFilterIndexScanRule.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/planner/physical/OpenSearchFilterIndexScanRule.java @@ -42,7 +42,8 @@ public void onMatch(RelOptRuleCall call) { } protected void apply(RelOptRuleCall call, Filter filter, CalciteOpenSearchIndexScan scan) { - if (scan.pushDownFilter(filter)) { + CalciteOpenSearchIndexScan newScan = scan.pushDownFilter(filter); + if (newScan != null) { call.transformTo(scan); } } diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java index 54a350c7428..c32a13a2305 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java @@ -7,6 +7,7 @@ import static java.util.Objects.requireNonNull; +import java.util.ArrayDeque; import java.util.List; import org.apache.calcite.adapter.enumerable.EnumerableRelImplementor; import org.apache.calcite.adapter.enumerable.PhysType; @@ -36,8 +37,6 @@ import org.opensearch.sql.opensearch.planner.physical.OpenSearchIndexRules; import org.opensearch.sql.opensearch.request.OpenSearchRequestBuilder; import org.opensearch.sql.opensearch.request.PredicateAnalyzer; -import org.opensearch.sql.opensearch.request.PredicateAnalyzer.ExpressionNotAnalyzableException; -import org.opensearch.sql.opensearch.request.PredicateAnalyzer.PredicateAnalyzerException; import org.opensearch.sql.opensearch.storage.OpenSearchIndex; /** Relational expression representing a scan of an OpenSearchIndex type. */ @@ -45,9 +44,12 @@ public class CalciteOpenSearchIndexScan extends OpenSearchTableScan { private static final Logger LOG = LogManager.getLogger(CalciteOpenSearchIndexScan.class); private final OpenSearchIndex osIndex; - private final OpenSearchRequestBuilder requestBuilder; + // The schema of this scan operator, it's initialized with the row type of the table, but may be + // changed by push down operations. private final RelDataType schema; + private final PushDownContext pushDownContext; + /** * Creates an CalciteOpenSearchIndexScan. * @@ -57,24 +59,31 @@ public class CalciteOpenSearchIndexScan extends OpenSearchTableScan { */ public CalciteOpenSearchIndexScan( RelOptCluster cluster, RelOptTable table, OpenSearchIndex index) { - this(cluster, table, index, index.createRequestBuilder(), table.getRowType()); + this(cluster, table, index, table.getRowType(), null); } - public CalciteOpenSearchIndexScan( + private CalciteOpenSearchIndexScan( RelOptCluster cluster, RelOptTable table, OpenSearchIndex index, - OpenSearchRequestBuilder requestBuilder, - RelDataType schema) { + RelDataType schema, + PushDownContext pushDownContext) { super(cluster, table); this.osIndex = requireNonNull(index, "OpenSearch index"); - this.requestBuilder = requestBuilder; this.schema = schema; + this.pushDownContext = pushDownContext == null ? new PushDownContext() : pushDownContext; + } + + public CalciteOpenSearchIndexScan copy() { + return new CalciteOpenSearchIndexScan( + getCluster(), table, osIndex, this.schema, pushDownContext.clone()); } public CalciteOpenSearchIndexScan copyWithNewSchema(RelDataType schema) { - // TODO: need to do deep-copy on requestBuilder in case non-idempotent push down. - return new CalciteOpenSearchIndexScan(getCluster(), table, osIndex, requestBuilder, schema); + // Do shallow copy for requestBuilder, thus requestBuilder among different plans produced in the + // optimization process won't affect each other. + return new CalciteOpenSearchIndexScan( + getCluster(), table, osIndex, schema, pushDownContext.clone()); } @Override @@ -115,6 +124,8 @@ public Result implement(EnumerableRelImplementor implementor, Prefer pref) { } public Enumerable<@Nullable Object> scan() { + OpenSearchRequestBuilder requestBuilder = osIndex.createRequestBuilder(); + pushDownContext.forEach(action -> action.apply(requestBuilder)); return new AbstractEnumerable<>() { @Override public Enumerator enumerator() { @@ -127,17 +138,18 @@ public Enumerator enumerator() { }; } - public boolean pushDownFilter(Filter filter) { + public CalciteOpenSearchIndexScan pushDownFilter(Filter filter) { try { + CalciteOpenSearchIndexScan newScan = this.copyWithNewSchema(filter.getRowType()); List schema = this.getRowType().getFieldNames(); QueryBuilder filterBuilder = PredicateAnalyzer.analyze(filter.getCondition(), schema); - requestBuilder.pushDownFilter(filterBuilder); + newScan.pushDownContext.add(requestBuilder -> requestBuilder.pushDownFilter(filterBuilder)); // TODO: handle the case where condition contains a score function - return true; - } catch (ExpressionNotAnalyzableException | PredicateAnalyzerException e) { + return newScan; + } catch (Exception e) { LOG.warn("Cannot analyze the filter condition {}", filter.getCondition(), e); } - return false; + return null; } /** @@ -152,7 +164,19 @@ public CalciteOpenSearchIndexScan pushDownProject(List selectedColumns) } RelDataType newSchema = builder.build(); CalciteOpenSearchIndexScan newScan = this.copyWithNewSchema(newSchema); - newScan.requestBuilder.pushDownProjectStream(newSchema.getFieldNames().stream()); + newScan.pushDownContext.add( + requestBuilder -> requestBuilder.pushDownProjectStream(newSchema.getFieldNames().stream())); return newScan; } + + static class PushDownContext extends ArrayDeque { + @Override + public PushDownContext clone() { + return (PushDownContext) super.clone(); + } + } + + private interface PushDownAction { + void apply(OpenSearchRequestBuilder requestBuilder); + } } From 6573ab7081fa08653dd88465ac32a910ecdce111 Mon Sep 17 00:00:00 2001 From: Heng Qian Date: Thu, 27 Feb 2025 19:30:36 +0800 Subject: [PATCH 5/5] Address comments and disable push down Signed-off-by: Heng Qian --- .../java/org/opensearch/sql/executor/QueryService.java | 3 +-- .../sql/calcite/standalone/CalcitePPLSortIT.java | 3 --- .../physical/OpenSearchFilterIndexScanRule.java | 2 +- .../sql/opensearch/setting/OpenSearchSettings.java | 8 +++++++- .../storage/scan/CalciteOpenSearchIndexScan.java | 10 +++++++--- 5 files changed, 16 insertions(+), 10 deletions(-) diff --git a/core/src/main/java/org/opensearch/sql/executor/QueryService.java b/core/src/main/java/org/opensearch/sql/executor/QueryService.java index ca67bd3efce..c2d581b8a7d 100644 --- a/core/src/main/java/org/opensearch/sql/executor/QueryService.java +++ b/core/src/main/java/org/opensearch/sql/executor/QueryService.java @@ -30,7 +30,6 @@ import org.apache.calcite.tools.Programs; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; -import org.jetbrains.annotations.NotNull; import org.opensearch.sql.analysis.AnalysisContext; import org.opensearch.sql.analysis.Analyzer; import org.opensearch.sql.ast.tree.UnresolvedPlan; @@ -146,7 +145,7 @@ public void executePlanByCalcite( * * @param osPlan Logical Plan derived from OpenSearch PPL */ - private static @NotNull RelNode convertToCalcitePlan(RelNode osPlan) { + private static RelNode convertToCalcitePlan(RelNode osPlan) { RelNode calcitePlan = osPlan; /* Calcite only ensures collation of the final result produced from the root sort operator. diff --git a/integ-test/src/test/java/org/opensearch/sql/calcite/standalone/CalcitePPLSortIT.java b/integ-test/src/test/java/org/opensearch/sql/calcite/standalone/CalcitePPLSortIT.java index 2100efaa3d8..545bed6e1a8 100644 --- a/integ-test/src/test/java/org/opensearch/sql/calcite/standalone/CalcitePPLSortIT.java +++ b/integ-test/src/test/java/org/opensearch/sql/calcite/standalone/CalcitePPLSortIT.java @@ -8,11 +8,8 @@ import static org.opensearch.sql.legacy.TestsConstants.TEST_INDEX_BANK; import java.io.IOException; -import org.junit.Ignore; import org.junit.jupiter.api.Test; -/** testSortXXAndXX could fail. TODO Remove this @Ignore when the issue fixed. */ -@Ignore public class CalcitePPLSortIT extends CalcitePPLIntegTestCase { @Override diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/planner/physical/OpenSearchFilterIndexScanRule.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/planner/physical/OpenSearchFilterIndexScanRule.java index 2b1e8e53124..621bfd8c6fd 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/planner/physical/OpenSearchFilterIndexScanRule.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/planner/physical/OpenSearchFilterIndexScanRule.java @@ -44,7 +44,7 @@ public void onMatch(RelOptRuleCall call) { protected void apply(RelOptRuleCall call, Filter filter, CalciteOpenSearchIndexScan scan) { CalciteOpenSearchIndexScan newScan = scan.pushDownFilter(filter); if (newScan != null) { - call.transformTo(scan); + call.transformTo(newScan); } } diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/setting/OpenSearchSettings.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/setting/OpenSearchSettings.java index f1f12a21553..96c74164025 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/setting/OpenSearchSettings.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/setting/OpenSearchSettings.java @@ -102,7 +102,7 @@ public class OpenSearchSettings extends Settings { public static final Setting CALCITE_PUSHDOWN_ENABLED_SETTING = Setting.boolSetting( Key.CALCITE_PUSHDOWN_ENABLED.getKeyValue(), - true, + false, Setting.Property.NodeScope, Setting.Property.Dynamic); @@ -309,6 +309,12 @@ public OpenSearchSettings(ClusterSettings clusterSettings) { Key.CALCITE_FALLBACK_ALLOWED, CALCITE_FALLBACK_ALLOWED_SETTING, new Updater(Key.CALCITE_FALLBACK_ALLOWED)); + register( + settingBuilder, + clusterSettings, + Key.CALCITE_PUSHDOWN_ENABLED, + CALCITE_PUSHDOWN_ENABLED_SETTING, + new Updater(Key.CALCITE_PUSHDOWN_ENABLED)); register( settingBuilder, clusterSettings, diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java index c32a13a2305..6c1386c5a3f 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteOpenSearchIndexScan.java @@ -47,7 +47,11 @@ public class CalciteOpenSearchIndexScan extends OpenSearchTableScan { // The schema of this scan operator, it's initialized with the row type of the table, but may be // changed by push down operations. private final RelDataType schema; - + // This context maintains all the push down actions, which will be applied to the requestBuilder + // when it begins to scan data from OpenSearch. + // Because OpenSearchRequestBuilder doesn't support deep copy while we want to keep the + // requestBuilder independent among different plans produced in the optimization process, + // so we cannot apply these actions right away. private final PushDownContext pushDownContext; /** @@ -59,7 +63,7 @@ public class CalciteOpenSearchIndexScan extends OpenSearchTableScan { */ public CalciteOpenSearchIndexScan( RelOptCluster cluster, RelOptTable table, OpenSearchIndex index) { - this(cluster, table, index, table.getRowType(), null); + this(cluster, table, index, table.getRowType(), new PushDownContext()); } private CalciteOpenSearchIndexScan( @@ -71,7 +75,7 @@ private CalciteOpenSearchIndexScan( super(cluster, table); this.osIndex = requireNonNull(index, "OpenSearch index"); this.schema = schema; - this.pushDownContext = pushDownContext == null ? new PushDownContext() : pushDownContext; + this.pushDownContext = pushDownContext; } public CalciteOpenSearchIndexScan copy() {