From fbf3ac23b6e1804c59f4e89cf9cfca16fdd4bca4 Mon Sep 17 00:00:00 2001 From: Heng Qian Date: Wed, 11 Feb 2026 11:21:30 +0800 Subject: [PATCH 1/4] Struct return array value instead of string Signed-off-by: Heng Qian --- .../sql/opensearch/executor/OpenSearchExecutionEngine.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/executor/OpenSearchExecutionEngine.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/executor/OpenSearchExecutionEngine.java index 58d797f4bf9..926fa4bc9b7 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/executor/OpenSearchExecutionEngine.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/executor/OpenSearchExecutionEngine.java @@ -230,7 +230,7 @@ public void execute( * Process values recursively, handling geo points and nested maps. Geo points are converted to * OpenSearchExprGeoPointValue. Maps are recursively processed to handle nested structures. */ - private static Object processValue(Object value) { + private static Object processValue(Object value) throws SQLException { if (value == null) { return null; } @@ -247,7 +247,7 @@ private static Object processValue(Object value) { return convertedMap; } if (value instanceof StructImpl) { - return ((StructImpl) value).toString(); + return List.of(((StructImpl) value).getAttributes()); } if (value instanceof List) { List list = (List) value; From 140acb9aeb0b258a29f33ab14f8470e4349f235d Mon Sep 17 00:00:00 2001 From: Heng Qian Date: Wed, 11 Feb 2026 13:38:59 +0800 Subject: [PATCH 2/4] Support filter in GraphLookup Signed-off-by: Heng Qian --- .../opensearch/sql/ast/tree/GraphLookup.java | 6 + .../sql/calcite/CalciteRelNodeVisitor.java | 12 +- .../sql/calcite/plan/rel/GraphLookup.java | 10 +- .../calcite/plan/rel/LogicalGraphLookup.java | 18 ++- docs/user/ppl/cmd/graphlookup.md | 24 +++- .../remote/CalcitePPLGraphLookupIT.java | 104 ++++++++++++++++++ .../rules/EnumerableGraphLookupRule.java | 3 +- .../scan/CalciteEnumerableGraphLookup.java | 31 +++++- ppl/src/main/antlr/OpenSearchPPLParser.g4 | 1 + .../opensearch/sql/ppl/parser/AstBuilder.java | 5 + .../sql/ppl/utils/PPLQueryDataAnonymizer.java | 6 + .../calcite/CalcitePPLGraphLookupTest.java | 36 ++++++ .../ppl/utils/PPLQueryDataAnonymizerTest.java | 14 +++ 13 files changed, 256 insertions(+), 14 deletions(-) diff --git a/core/src/main/java/org/opensearch/sql/ast/tree/GraphLookup.java b/core/src/main/java/org/opensearch/sql/ast/tree/GraphLookup.java index 51084d7b9b9..6771c35c43d 100644 --- a/core/src/main/java/org/opensearch/sql/ast/tree/GraphLookup.java +++ b/core/src/main/java/org/opensearch/sql/ast/tree/GraphLookup.java @@ -18,6 +18,7 @@ import org.opensearch.sql.ast.AbstractNodeVisitor; import org.opensearch.sql.ast.expression.Field; import org.opensearch.sql.ast.expression.Literal; +import org.opensearch.sql.ast.expression.UnresolvedExpression; /** * AST node for graphLookup command. Performs BFS graph traversal on a lookup table. @@ -74,6 +75,11 @@ public enum Direction { /** Whether to use PIT (Point In Time) search for the lookup table to get complete results. */ private final boolean usePIT; + /** + * Optional filter condition to restrict which lookup table documents participate in traversal. + */ + private @Nullable final UnresolvedExpression filter; + private UnresolvedPlan child; public String getDepthFieldName() { diff --git a/core/src/main/java/org/opensearch/sql/calcite/CalciteRelNodeVisitor.java b/core/src/main/java/org/opensearch/sql/calcite/CalciteRelNodeVisitor.java index c9227781f6a..9a16f2bf802 100644 --- a/core/src/main/java/org/opensearch/sql/calcite/CalciteRelNodeVisitor.java +++ b/core/src/main/java/org/opensearch/sql/calcite/CalciteRelNodeVisitor.java @@ -2606,9 +2606,16 @@ public RelNode visitGraphLookup(GraphLookup node, CalcitePlanContext context) { // 3. Visit and materialize lookup table analyze(node.getFromTable(), context); tryToRemoveMetaFields(context, true); + + // 4. Convert filter expression to RexNode against lookup table schema + RexNode filterRex = null; + if (node.getFilter() != null) { + filterRex = rexVisitor.analyze(node.getFilter(), context); + } + RelNode lookupTable = builder.build(); - // 4. Create LogicalGraphLookup RelNode + // 5. Create LogicalGraphLookup RelNode // The conversion rule will extract the OpenSearchIndex from the lookup table RelNode graphLookup = LogicalGraphLookup.create( @@ -2623,7 +2630,8 @@ public RelNode visitGraphLookup(GraphLookup node, CalcitePlanContext context) { bidirectional, supportArray, batchMode, - usePIT); + usePIT, + filterRex); builder.push(graphLookup); return builder.peek(); diff --git a/core/src/main/java/org/opensearch/sql/calcite/plan/rel/GraphLookup.java b/core/src/main/java/org/opensearch/sql/calcite/plan/rel/GraphLookup.java index ef7134a0162..02ed97faf0c 100644 --- a/core/src/main/java/org/opensearch/sql/calcite/plan/rel/GraphLookup.java +++ b/core/src/main/java/org/opensearch/sql/calcite/plan/rel/GraphLookup.java @@ -16,6 +16,7 @@ import org.apache.calcite.rel.metadata.RelMetadataQuery; import org.apache.calcite.rel.type.RelDataType; import org.apache.calcite.rel.type.RelDataTypeFactory; +import org.apache.calcite.rex.RexNode; import org.apache.calcite.sql.type.SqlTypeName; /** @@ -51,6 +52,7 @@ public abstract class GraphLookup extends BiRel { protected final boolean supportArray; protected final boolean batchMode; protected final boolean usePIT; + @Nullable protected final RexNode filter; private RelDataType outputRowType; @@ -72,6 +74,7 @@ public abstract class GraphLookup extends BiRel { * pushdown) * @param batchMode Whether to batch all source start values into a single unified BFS * @param usePIT Whether to use PIT (Point In Time) search for complete results + * @param filter Optional filter condition for lookup table documents */ protected GraphLookup( RelOptCluster cluster, @@ -87,7 +90,8 @@ protected GraphLookup( boolean bidirectional, boolean supportArray, boolean batchMode, - boolean usePIT) { + boolean usePIT, + @Nullable RexNode filter) { super(cluster, traitSet, source, lookup); this.startField = startField; this.fromField = fromField; @@ -99,6 +103,7 @@ protected GraphLookup( this.supportArray = supportArray; this.batchMode = batchMode; this.usePIT = usePIT; + this.filter = filter; } /** Returns the source table RelNode. */ @@ -181,6 +186,7 @@ public RelWriter explainTerms(RelWriter pw) { .item("bidirectional", bidirectional) .itemIf("supportArray", supportArray, supportArray) .itemIf("batchMode", batchMode, batchMode) - .itemIf("usePIT", usePIT, usePIT); + .itemIf("usePIT", usePIT, usePIT) + .itemIf("filter", filter, filter != null); } } diff --git a/core/src/main/java/org/opensearch/sql/calcite/plan/rel/LogicalGraphLookup.java b/core/src/main/java/org/opensearch/sql/calcite/plan/rel/LogicalGraphLookup.java index b02bbec3742..94db3689f8c 100644 --- a/core/src/main/java/org/opensearch/sql/calcite/plan/rel/LogicalGraphLookup.java +++ b/core/src/main/java/org/opensearch/sql/calcite/plan/rel/LogicalGraphLookup.java @@ -12,6 +12,7 @@ import org.apache.calcite.plan.RelOptCluster; import org.apache.calcite.plan.RelTraitSet; import org.apache.calcite.rel.RelNode; +import org.apache.calcite.rex.RexNode; /** * Logical RelNode for graphLookup command. TODO: need to support trim fields and several transpose @@ -37,6 +38,7 @@ public class LogicalGraphLookup extends GraphLookup { * @param supportArray Whether to support array-typed fields * @param batchMode Whether to batch all source start values into a single unified BFS * @param usePIT Whether to use PIT (Point In Time) search for complete results + * @param filter Optional filter condition for lookup table documents */ protected LogicalGraphLookup( RelOptCluster cluster, @@ -52,7 +54,8 @@ protected LogicalGraphLookup( boolean bidirectional, boolean supportArray, boolean batchMode, - boolean usePIT) { + boolean usePIT, + @Nullable RexNode filter) { super( cluster, traitSet, @@ -67,7 +70,8 @@ protected LogicalGraphLookup( bidirectional, supportArray, batchMode, - usePIT); + usePIT, + filter); } /** @@ -85,6 +89,7 @@ protected LogicalGraphLookup( * @param supportArray Whether to support array-typed fields * @param batchMode Whether to batch all source start values into a single unified BFS * @param usePIT Whether to use PIT (Point In Time) search for complete results + * @param filter Optional filter condition for lookup table documents * @return A new LogicalGraphLookup instance */ public static LogicalGraphLookup create( @@ -99,7 +104,8 @@ public static LogicalGraphLookup create( boolean bidirectional, boolean supportArray, boolean batchMode, - boolean usePIT) { + boolean usePIT, + @Nullable RexNode filter) { RelOptCluster cluster = source.getCluster(); RelTraitSet traitSet = cluster.traitSetOf(Convention.NONE); return new LogicalGraphLookup( @@ -116,7 +122,8 @@ public static LogicalGraphLookup create( bidirectional, supportArray, batchMode, - usePIT); + usePIT, + filter); } @Override @@ -135,6 +142,7 @@ public RelNode copy(RelTraitSet traitSet, List inputs) { bidirectional, supportArray, batchMode, - usePIT); + usePIT, + filter); } } diff --git a/docs/user/ppl/cmd/graphlookup.md b/docs/user/ppl/cmd/graphlookup.md index e768e02f6b8..0bb8ba3d196 100644 --- a/docs/user/ppl/cmd/graphlookup.md +++ b/docs/user/ppl/cmd/graphlookup.md @@ -8,7 +8,7 @@ The `graphLookup` command performs recursive graph traversal on a collection usi The `graphLookup` command has the following syntax: ```syntax -graphLookup startField= fromField= toField= [maxDepth=] [depthField=] [direction=(uni | bi)] [supportArray=(true | false)] [batchMode=(true | false)] [usePIT=(true | false)] as +graphLookup startField= fromField= toField= [maxDepth=] [depthField=] [direction=(uni | bi)] [supportArray=(true | false)] [batchMode=(true | false)] [usePIT=(true | false)] [filter=()] as ``` The following are examples of the `graphLookup` command syntax: @@ -20,6 +20,7 @@ source = employees | graphLookup employees startField=reportsTo fromField=report source = employees | graphLookup employees startField=reportsTo fromField=reportsTo toField=name direction=bi as connections source = travelers | graphLookup airports startField=nearestAirport fromField=connects toField=airport supportArray=true as reachableAirports source = airports | graphLookup airports startField=airport fromField=connects toField=airport supportArray=true as reachableAirports +source = employees | graphLookup employees startField=reportsTo fromField=reportsTo toField=name filter=(status = 'active' AND age > 18) as reportingHierarchy ``` ## Parameters @@ -38,6 +39,7 @@ The `graphLookup` command supports the following parameters. | `supportArray=(true \| false)` | Optional | When `true`, disables early visited-node filter pushdown to OpenSearch. Default is `false`. Set to `true` when `fromField` or `toField` contains array values to ensure correct traversal behavior. See [Array Field Handling](#array-field-handling) for details. | | `batchMode=(true \| false)` | Optional | When `true`, collects all start values from all source rows and performs a single unified BFS traversal. Default is `false`. The output changes to two arrays: `[Array, Array]`. See [Batch Mode](#batch-mode) for details. | | `usePIT=(true \| false)` | Optional | When `true`, enables PIT (Point In Time) search for the lookup table, allowing paginated retrieval of complete results without the `max_result_window` size limit. Default is `false`. See [PIT Search](#pit-search) for details. | +| `filter=()` | Optional | A filter condition to restrict which lookup table documents participate in the graph traversal. Only documents matching the condition are considered as candidates during BFS. Parentheses around the condition are required. Example: `filter=(status = 'active' AND age > 18)`. | | `as ` | Required | The name of the output array field that will contain all documents found during the graph traversal. | ## How It Works @@ -329,6 +331,26 @@ source = employees as reportingHierarchy ``` +## Filtered Graph Traversal + +The `filter` parameter restricts which documents in the lookup table are considered during the BFS traversal. Only documents matching the filter condition participate as candidates at each traversal level. + +### Example + +The following query traverses only active employees in the reporting hierarchy: + +```ppl ignore +source = employees + | graphLookup employees + startField=reportsTo + fromField=reportsTo + toField=name + filter=(status = 'active') + as reportingHierarchy +``` + +The filter is applied at the OpenSearch query level, so it combines efficiently with the BFS traversal queries. At each BFS level, the query sent to OpenSearch is effectively: `bool { filter: [user_filter, bfs_terms_query] }`. + ## Limitations - The source input, which provides the starting point for the traversal, has a limitation of 100 documents to avoid performance issues. diff --git a/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePPLGraphLookupIT.java b/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePPLGraphLookupIT.java index 498b17dab91..abe84b34abc 100644 --- a/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePPLGraphLookupIT.java +++ b/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePPLGraphLookupIT.java @@ -395,6 +395,110 @@ public void testBidirectionalAirportConnections() throws IOException { result, rows("ORD", List.of("JFK"), List.of("{JFK, [BOS, ORD]}", "{BOS, [JFK, PWM]}"))); } + // ==================== Filter Tests ==================== + + /** + * Test: Filter employee hierarchy by id. Only lookup documents with id > 3 (Andrew=4, Asya=5, + * Dan=6) participate in traversal. Dev starts with reportsTo=Eliot, but Eliot (id=2) is excluded + * by filter, so Dev gets empty results. + */ + @Test + public void testEmployeeHierarchyWithFilter() throws IOException { + JSONObject result = + executeQuery( + String.format( + "source=%s" + + " | graphLookup %s" + + " startField=reportsTo" + + " fromField=reportsTo" + + " toField=name" + + " filter=(id > 3)" + + " as reportingHierarchy", + TEST_INDEX_GRAPH_EMPLOYEES, TEST_INDEX_GRAPH_EMPLOYEES)); + + verifySchema( + result, + schema("name", "string"), + schema("reportsTo", "string"), + schema("id", "int"), + schema("reportingHierarchy", "array")); + // Only documents with id > 3 (Andrew=4, Asya=5, Dan=6) are in lookup table + // Dev: reportsTo=Eliot -> Eliot(id=2) is filtered out -> empty + // Eliot: reportsTo=Ron -> Ron(id=3) is filtered out -> empty + // Ron: reportsTo=Andrew -> Andrew(id=4) passes filter -> [{Andrew, null, 4}] + // Andrew: reportsTo=null -> empty + // Asya: reportsTo=Ron -> Ron(id=3) is filtered out -> empty + // Dan: reportsTo=Andrew -> Andrew(id=4) passes filter -> [{Andrew, null, 4}] + verifyDataRows( + result, + rows("Dev", "Eliot", 1, Collections.emptyList()), + rows("Eliot", "Ron", 2, Collections.emptyList()), + rows("Ron", "Andrew", 3, List.of("{Andrew, null, 4}")), + rows("Andrew", null, 4, Collections.emptyList()), + rows("Asya", "Ron", 5, Collections.emptyList()), + rows("Dan", "Andrew", 6, List.of("{Andrew, null, 4}"))); + } + + /** + * Test: Filter employee hierarchy with keyword match. Only employees whose name is NOT 'Andrew' + * participate in traversal. + */ + @Test + public void testEmployeeHierarchyWithKeywordFilter() throws IOException { + JSONObject result = + executeQuery( + String.format( + "source=%s" + + " | where name = 'Ron'" + + " | graphLookup %s" + + " startField=reportsTo" + + " fromField=reportsTo" + + " toField=name" + + " filter=(name != 'Andrew')" + + " as reportingHierarchy", + TEST_INDEX_GRAPH_EMPLOYEES, TEST_INDEX_GRAPH_EMPLOYEES)); + + verifySchema( + result, + schema("name", "string"), + schema("reportsTo", "string"), + schema("id", "int"), + schema("reportingHierarchy", "array")); + // Ron: reportsTo=Andrew -> Andrew is filtered out by name != 'Andrew' -> empty + verifyDataRows(result, rows("Ron", "Andrew", 3, Collections.emptyList())); + } + + /** + * Test: Filter with maxDepth combined. Dev traverses reporting chain but only considers lookup + * documents with id <= 3. + */ + @Test + public void testEmployeeHierarchyWithFilterAndMaxDepth() throws IOException { + JSONObject result = + executeQuery( + String.format( + "source=%s" + + " | where name = 'Dev'" + + " | graphLookup %s" + + " startField=reportsTo" + + " fromField=reportsTo" + + " toField=name" + + " maxDepth=3" + + " filter=(id <= 3)" + + " as reportingHierarchy", + TEST_INDEX_GRAPH_EMPLOYEES, TEST_INDEX_GRAPH_EMPLOYEES)); + + verifySchema( + result, + schema("name", "string"), + schema("reportsTo", "string"), + schema("id", "int"), + schema("reportingHierarchy", "array")); + // Dev: reportsTo=Eliot -> Eliot(id=2) passes -> then Eliot.reportsTo=Ron -> Ron(id=3) passes + // -> then Ron.reportsTo=Andrew -> Andrew(id=4) is filtered out -> stops + verifyDataRows(result, rows("Dev", "Eliot", 1, List.of("{Eliot, Ron, 2}", "{Ron, Andrew, 3}"))); + } + // ==================== Edge Cases ==================== /** Test 13: Graph lookup on empty result set (non-existent employee). */ diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/planner/rules/EnumerableGraphLookupRule.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/planner/rules/EnumerableGraphLookupRule.java index f76c90ab47d..e210095b480 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/planner/rules/EnumerableGraphLookupRule.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/planner/rules/EnumerableGraphLookupRule.java @@ -101,6 +101,7 @@ public RelNode convert(RelNode rel) { graphLookup.isBidirectional(), graphLookup.isSupportArray(), graphLookup.isBatchMode(), - graphLookup.isUsePIT()); + graphLookup.isUsePIT(), + graphLookup.getFilter()); } } diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteEnumerableGraphLookup.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteEnumerableGraphLookup.java index 2cd0f4a746f..6307a741468 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteEnumerableGraphLookup.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/storage/scan/CalciteEnumerableGraphLookup.java @@ -14,6 +14,7 @@ import java.util.HashSet; import java.util.Iterator; import java.util.List; +import java.util.Map; import java.util.Queue; import java.util.Set; import lombok.Getter; @@ -32,6 +33,7 @@ import org.apache.calcite.plan.RelTraitSet; import org.apache.calcite.rel.RelNode; import org.apache.calcite.rel.metadata.RelMetadataQuery; +import org.apache.calcite.rex.RexNode; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import org.checkerframework.checker.nullness.qual.Nullable; @@ -39,6 +41,8 @@ import org.opensearch.index.query.QueryBuilders; import org.opensearch.sql.calcite.plan.Scannable; import org.opensearch.sql.calcite.plan.rel.GraphLookup; +import org.opensearch.sql.data.type.ExprType; +import org.opensearch.sql.opensearch.request.PredicateAnalyzer; import org.opensearch.sql.opensearch.request.PredicateAnalyzer.NamedFieldExpression; import org.opensearch.sql.opensearch.storage.scan.context.LimitDigest; import org.opensearch.sql.opensearch.storage.scan.context.OSRequestBuilderAction; @@ -74,6 +78,7 @@ public class CalciteEnumerableGraphLookup extends GraphLookup implements Enumera * @param supportArray Whether to support array-typed fields * @param batchMode Whether to batch all source start values into a single unified BFS * @param usePIT Whether to use PIT (Point In Time) search for complete results + * @param filter Optional filter condition for lookup table documents */ public CalciteEnumerableGraphLookup( RelOptCluster cluster, @@ -89,7 +94,8 @@ public CalciteEnumerableGraphLookup( boolean bidirectional, boolean supportArray, boolean batchMode, - boolean usePIT) { + boolean usePIT, + @Nullable RexNode filter) { super( cluster, traitSet, @@ -104,7 +110,8 @@ public CalciteEnumerableGraphLookup( bidirectional, supportArray, batchMode, - usePIT); + usePIT, + filter); } @Override @@ -123,7 +130,8 @@ public RelNode copy(RelTraitSet traitSet, List inputs) { bidirectional, supportArray, batchMode, - usePIT); + usePIT, + filter); } @Override @@ -209,6 +217,23 @@ private static class GraphLookupEnumerator implements Enumerator<@Nullable Objec this.startFieldIndex = sourceFields.indexOf(graphLookup.getStartField()); this.fromFieldIdx = lookupFields.indexOf(graphLookup.fromField); this.toFieldIdx = lookupFields.indexOf(graphLookup.toField); + + // Push down user-specified filter to the lookup scan + if (graphLookup.filter != null) { + List schema = graphLookup.getLookup().getRowType().getFieldNames(); + Map fieldTypes = this.lookupScan.getOsIndex().getAllFieldTypes(); + try { + QueryBuilder filterQuery = + PredicateAnalyzer.analyze(graphLookup.filter, schema, fieldTypes); + this.lookupScan.pushDownContext.add( + PushDownType.FILTER, + null, + (OSRequestBuilderAction) rb -> rb.pushDownFilterForCalcite(filterQuery)); + } catch (PredicateAnalyzer.ExpressionNotAnalyzableException e) { + throw new RuntimeException( + "Cannot push down filter for graphLookup: " + e.getMessage(), e); + } + } } @Override diff --git a/ppl/src/main/antlr/OpenSearchPPLParser.g4 b/ppl/src/main/antlr/OpenSearchPPLParser.g4 index d860874a1e9..7b9e1a0b5f4 100644 --- a/ppl/src/main/antlr/OpenSearchPPLParser.g4 +++ b/ppl/src/main/antlr/OpenSearchPPLParser.g4 @@ -641,6 +641,7 @@ graphLookupOption | (SUPPORT_ARRAY EQUAL booleanLiteral) | (BATCH_MODE EQUAL booleanLiteral) | (USE_PIT EQUAL booleanLiteral) + | (FILTER EQUAL LT_PRTHS logicalExpression RT_PRTHS) ; // clauses diff --git a/ppl/src/main/java/org/opensearch/sql/ppl/parser/AstBuilder.java b/ppl/src/main/java/org/opensearch/sql/ppl/parser/AstBuilder.java index 721f13033f3..9c367d73985 100644 --- a/ppl/src/main/java/org/opensearch/sql/ppl/parser/AstBuilder.java +++ b/ppl/src/main/java/org/opensearch/sql/ppl/parser/AstBuilder.java @@ -1499,6 +1499,7 @@ public UnresolvedPlan visitGraphLookupCommand(OpenSearchPPLParser.GraphLookupCom boolean supportArray = false; boolean batchMode = false; boolean usePIT = false; + UnresolvedExpression filter = null; for (OpenSearchPPLParser.GraphLookupOptionContext option : ctx.graphLookupOption()) { if (option.FROM_FIELD() != null) { @@ -1531,6 +1532,9 @@ public UnresolvedPlan visitGraphLookupCommand(OpenSearchPPLParser.GraphLookupCom Literal literal = (Literal) internalVisitExpression(option.booleanLiteral()); usePIT = Boolean.TRUE.equals(literal.getValue()); } + if (option.FILTER() != null) { + filter = internalVisitExpression(option.logicalExpression()); + } } Field as = (Field) internalVisitExpression(ctx.outputField); @@ -1551,6 +1555,7 @@ public UnresolvedPlan visitGraphLookupCommand(OpenSearchPPLParser.GraphLookupCom .supportArray(supportArray) .batchMode(batchMode) .usePIT(usePIT) + .filter(filter) .build(); } } diff --git a/ppl/src/main/java/org/opensearch/sql/ppl/utils/PPLQueryDataAnonymizer.java b/ppl/src/main/java/org/opensearch/sql/ppl/utils/PPLQueryDataAnonymizer.java index 2402fac8700..90f4ce92724 100644 --- a/ppl/src/main/java/org/opensearch/sql/ppl/utils/PPLQueryDataAnonymizer.java +++ b/ppl/src/main/java/org/opensearch/sql/ppl/utils/PPLQueryDataAnonymizer.java @@ -251,6 +251,12 @@ public String visitGraphLookup(GraphLookup node, String context) { if (node.isUsePIT()) { command.append(" usePIT=true"); } + if (node.getFilter() != null) { + command + .append(" filter=(") + .append(expressionAnalyzer.analyze(node.getFilter(), context)) + .append(")"); + } command.append(" as ").append(MASK_COLUMN); return command.toString(); } diff --git a/ppl/src/test/java/org/opensearch/sql/ppl/calcite/CalcitePPLGraphLookupTest.java b/ppl/src/test/java/org/opensearch/sql/ppl/calcite/CalcitePPLGraphLookupTest.java index da169704925..2790f3aadd5 100644 --- a/ppl/src/test/java/org/opensearch/sql/ppl/calcite/CalcitePPLGraphLookupTest.java +++ b/ppl/src/test/java/org/opensearch/sql/ppl/calcite/CalcitePPLGraphLookupTest.java @@ -93,6 +93,42 @@ public void testGraphLookupWithMaxDepth() { verifyLogical(root, expectedLogical); } + @Test + public void testGraphLookupWithFilter() { + // Test graphLookup with filter parameter + String ppl = + "source=employee | graphLookup employee startField=reportsTo fromField=reportsTo" + + " toField=name filter=(id > 2) as reportingHierarchy"; + + RelNode root = getRelNode(ppl); + String expectedLogical = + "LogicalGraphLookup(fromField=[reportsTo], toField=[name]," + + " outputField=[reportingHierarchy], depthField=[null], maxDepth=[0]," + + " bidirectional=[false], filter=[>($0, 2)])\n" + + " LogicalSort(fetch=[100])\n" + + " LogicalTableScan(table=[[scott, employee]])\n" + + " LogicalTableScan(table=[[scott, employee]])\n"; + verifyLogical(root, expectedLogical); + } + + @Test + public void testGraphLookupWithCompoundFilter() { + // Test graphLookup with compound filter condition + String ppl = + "source=employee | graphLookup employee startField=reportsTo fromField=reportsTo" + + " toField=name filter=(id > 1 AND name != 'Andrew') as reportingHierarchy"; + + RelNode root = getRelNode(ppl); + String expectedLogical = + "LogicalGraphLookup(fromField=[reportsTo], toField=[name]," + + " outputField=[reportingHierarchy], depthField=[null], maxDepth=[0]," + + " bidirectional=[false], filter=[AND(>($0, 1), <>($1, 'Andrew'))])\n" + + " LogicalSort(fetch=[100])\n" + + " LogicalTableScan(table=[[scott, employee]])\n" + + " LogicalTableScan(table=[[scott, employee]])\n"; + verifyLogical(root, expectedLogical); + } + @Test public void testGraphLookupBidirectional() { // Test graphLookup with bidirectional traversal diff --git a/ppl/src/test/java/org/opensearch/sql/ppl/utils/PPLQueryDataAnonymizerTest.java b/ppl/src/test/java/org/opensearch/sql/ppl/utils/PPLQueryDataAnonymizerTest.java index 943db6c2ba2..0398d30bf17 100644 --- a/ppl/src/test/java/org/opensearch/sql/ppl/utils/PPLQueryDataAnonymizerTest.java +++ b/ppl/src/test/java/org/opensearch/sql/ppl/utils/PPLQueryDataAnonymizerTest.java @@ -695,6 +695,20 @@ public void testGraphLookup() { anonymize( "source=t | graphLookup employees fromField=manager toField=name" + " batchMode=true as reportingHierarchy")); + // graphLookup with filter + assertEquals( + "source=table | graphlookup table fromField=identifier toField=identifier" + + " direction=uni filter=(identifier = ***) as identifier", + anonymize( + "source=t | graphLookup employees fromField=manager toField=name" + + " filter=(status = 'active') as reportingHierarchy")); + // graphLookup with compound filter + assertEquals( + "source=table | graphlookup table fromField=identifier toField=identifier" + + " direction=uni filter=(identifier = *** and identifier > ***) as identifier", + anonymize( + "source=t | graphLookup employees fromField=manager toField=name" + + " filter=(status = 'active' AND id > 2) as reportingHierarchy")); } @Test From 407e256d40cd795e07bed7d8f370aa72ab31d230 Mon Sep 17 00:00:00 2001 From: Heng Qian Date: Wed, 11 Feb 2026 15:03:01 +0800 Subject: [PATCH 3/4] Fix IT Signed-off-by: Heng Qian --- .../remote/CalcitePPLGraphLookupIT.java | 127 +++++++++++------- .../executor/OpenSearchExecutionEngine.java | 3 +- 2 files changed, 81 insertions(+), 49 deletions(-) diff --git a/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePPLGraphLookupIT.java b/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePPLGraphLookupIT.java index abe84b34abc..fbaefdb8c3c 100644 --- a/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePPLGraphLookupIT.java +++ b/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePPLGraphLookupIT.java @@ -14,6 +14,7 @@ import static org.opensearch.sql.util.MatcherUtils.verifySchema; import java.io.IOException; +import java.util.Arrays; import java.util.Collections; import java.util.List; import org.json.JSONObject; @@ -72,12 +73,12 @@ public void testEmployeeHierarchyBasicTraversal() throws IOException { schema("reportingHierarchy", "array")); verifyDataRows( result, - rows("Dev", "Eliot", 1, List.of("{Eliot, Ron, 2}")), - rows("Eliot", "Ron", 2, List.of("{Ron, Andrew, 3}")), - rows("Ron", "Andrew", 3, List.of("{Andrew, null, 4}")), + rows("Dev", "Eliot", 1, List.of(List.of("Eliot", "Ron", 2))), + rows("Eliot", "Ron", 2, List.of(List.of("Ron", "Andrew", 3))), + rows("Ron", "Andrew", 3, List.of(Arrays.asList("Andrew", null, 4))), rows("Andrew", null, 4, Collections.emptyList()), - rows("Asya", "Ron", 5, List.of("{Ron, Andrew, 3}")), - rows("Dan", "Andrew", 6, List.of("{Andrew, null, 4}"))); + rows("Asya", "Ron", 5, List.of(List.of("Ron", "Andrew", 3))), + rows("Dan", "Andrew", 6, List.of(Arrays.asList("Andrew", null, 4)))); } /** Test 2: Employee hierarchy traversal with depth field. */ @@ -103,12 +104,12 @@ public void testEmployeeHierarchyWithDepthField() throws IOException { schema("reportingHierarchy", "array")); verifyDataRows( result, - rows("Dev", "Eliot", 1, List.of("{Eliot, Ron, 2, 0}")), - rows("Eliot", "Ron", 2, List.of("{Ron, Andrew, 3, 0}")), - rows("Ron", "Andrew", 3, List.of("{Andrew, null, 4, 0}")), + rows("Dev", "Eliot", 1, List.of(List.of("Eliot", "Ron", 2, 0))), + rows("Eliot", "Ron", 2, List.of(List.of("Ron", "Andrew", 3, 0))), + rows("Ron", "Andrew", 3, List.of(Arrays.asList("Andrew", null, 4, 0))), rows("Andrew", null, 4, Collections.emptyList()), - rows("Asya", "Ron", 5, List.of("{Ron, Andrew, 3, 0}")), - rows("Dan", "Andrew", 6, List.of("{Andrew, null, 4, 0}"))); + rows("Asya", "Ron", 5, List.of(List.of("Ron", "Andrew", 3, 0))), + rows("Dan", "Andrew", 6, List.of(Arrays.asList("Andrew", null, 4, 0)))); } /** Test 3: Employee hierarchy with maxDepth=1 (allows 2 levels of traversal). */ @@ -134,12 +135,20 @@ public void testEmployeeHierarchyWithMaxDepth() throws IOException { schema("reportingHierarchy", "array")); verifyDataRows( result, - rows("Dev", "Eliot", 1, List.of("{Eliot, Ron, 2}", "{Ron, Andrew, 3}")), - rows("Eliot", "Ron", 2, List.of("{Ron, Andrew, 3}", "{Andrew, null, 4}")), - rows("Ron", "Andrew", 3, List.of("{Andrew, null, 4}")), + rows("Dev", "Eliot", 1, List.of(List.of("Eliot", "Ron", 2), List.of("Ron", "Andrew", 3))), + rows( + "Eliot", + "Ron", + 2, + List.of(List.of("Ron", "Andrew", 3), Arrays.asList("Andrew", null, 4))), + rows("Ron", "Andrew", 3, List.of(Arrays.asList("Andrew", null, 4))), rows("Andrew", null, 4, Collections.emptyList()), - rows("Asya", "Ron", 5, List.of("{Ron, Andrew, 3}", "{Andrew, null, 4}")), - rows("Dan", "Andrew", 6, List.of("{Andrew, null, 4}"))); + rows( + "Asya", + "Ron", + 5, + List.of(List.of("Ron", "Andrew", 3), Arrays.asList("Andrew", null, 4))), + rows("Dan", "Andrew", 6, List.of(Arrays.asList("Andrew", null, 4)))); } /** Test 4: Query Dev's complete reporting chain: Dev->Eliot->Ron->Andrew */ @@ -163,7 +172,7 @@ public void testEmployeeHierarchyForSpecificEmployee() throws IOException { schema("reportsTo", "string"), schema("id", "int"), schema("reportingHierarchy", "array")); - verifyDataRows(result, rows("Dev", "Eliot", 1, List.of("{Eliot, Ron, 2}"))); + verifyDataRows(result, rows("Dev", "Eliot", 1, List.of(List.of("Eliot", "Ron", 2)))); } // ==================== Airport Connections Tests ==================== @@ -190,11 +199,11 @@ public void testAirportConnections() throws IOException { schema("reachableAirports", "array")); verifyDataRows( result, - rows("JFK", List.of("BOS", "ORD"), List.of("{JFK, [BOS, ORD]}")), - rows("BOS", List.of("JFK", "PWM"), List.of("{BOS, [JFK, PWM]}")), - rows("ORD", List.of("JFK"), List.of("{ORD, [JFK]}")), - rows("PWM", List.of("BOS", "LHR"), List.of("{PWM, [BOS, LHR]}")), - rows("LHR", List.of("PWM"), List.of("{LHR, [PWM]}"))); + rows("JFK", List.of("BOS", "ORD"), List.of(List.of("JFK", List.of("BOS", "ORD")))), + rows("BOS", List.of("JFK", "PWM"), List.of(List.of("BOS", List.of("JFK", "PWM")))), + rows("ORD", List.of("JFK"), List.of(List.of("ORD", List.of("JFK")))), + rows("PWM", List.of("BOS", "LHR"), List.of(List.of("PWM", List.of("BOS", "LHR")))), + rows("LHR", List.of("PWM"), List.of(List.of("LHR", List.of("PWM"))))); } /** Test 6: Find airports reachable from JFK within maxDepth=1. */ @@ -221,7 +230,10 @@ public void testAirportConnectionsWithMaxDepth() throws IOException { schema("reachableAirports", "array")); verifyDataRows( result, - rows("JFK", List.of("BOS", "ORD"), List.of("{JFK, [BOS, ORD]}", "{BOS, [JFK, PWM]}"))); + rows( + "JFK", + List.of("BOS", "ORD"), + List.of(List.of("JFK", List.of("BOS", "ORD")), List.of("BOS", List.of("JFK", "PWM"))))); } /** Test 7: Find airports with default depth(=0) and start value of list */ @@ -244,7 +256,9 @@ public void testAirportConnectionsWithDepthField() throws IOException { schema("airport", "string"), schema("connects", "string"), schema("reachableAirports", "array")); - verifyDataRows(result, rows("JFK", List.of("BOS", "ORD"), List.of("{BOS, [JFK, PWM], 0}"))); + verifyDataRows( + result, + rows("JFK", List.of("BOS", "ORD"), List.of(List.of("BOS", List.of("JFK", "PWM"), 0)))); } /** @@ -271,9 +285,9 @@ public void testTravelersReachableAirports() throws IOException { schema("reachableAirports", "array")); verifyDataRows( result, - rows("Dev", "JFK", List.of("{JFK, [BOS, ORD]}")), - rows("Eliot", "JFK", List.of("{JFK, [BOS, ORD]}")), - rows("Jeff", "BOS", List.of("{BOS, [JFK, PWM]}"))); + rows("Dev", "JFK", List.of(List.of("JFK", List.of("BOS", "ORD")))), + rows("Eliot", "JFK", List.of(List.of("JFK", List.of("BOS", "ORD")))), + rows("Jeff", "BOS", List.of(List.of("BOS", List.of("JFK", "PWM"))))); } /** @@ -300,7 +314,7 @@ public void testTravelerReachableAirportsWithDepthField() throws IOException { schema("name", "string"), schema("nearestAirport", "string"), schema("reachableAirports", "array")); - verifyDataRows(result, rows("Dev", "JFK", List.of("{JFK, [BOS, ORD], 0}"))); + verifyDataRows(result, rows("Dev", "JFK", List.of(List.of("JFK", List.of("BOS", "ORD"), 0)))); } /** @@ -331,7 +345,12 @@ public void testTravelerReachableAirportsWithMaxDepth() throws IOException { verifyDataRows( result, rows( - "Jeff", "BOS", List.of("{BOS, [JFK, PWM]}", "{JFK, [BOS, ORD]}", "{PWM, [BOS, LHR]}"))); + "Jeff", + "BOS", + List.of( + List.of("BOS", List.of("JFK", "PWM")), + List.of("JFK", List.of("BOS", "ORD")), + List.of("PWM", List.of("BOS", "LHR"))))); } // ==================== Bidirectional Traversal Tests ==================== @@ -364,7 +383,10 @@ public void testBidirectionalEmployeeHierarchy() throws IOException { "Ron", "Andrew", 3, - List.of("{Ron, Andrew, 3}", "{Andrew, null, 4}", "{Dan, Andrew, 6}"))); + List.of( + List.of("Ron", "Andrew", 3), + Arrays.asList("Andrew", null, 4), + List.of("Dan", "Andrew", 6)))); } /** @@ -392,7 +414,11 @@ public void testBidirectionalAirportConnections() throws IOException { schema("connects", "string"), schema("allConnections", "array")); verifyDataRows( - result, rows("ORD", List.of("JFK"), List.of("{JFK, [BOS, ORD]}", "{BOS, [JFK, PWM]}"))); + result, + rows( + "ORD", + List.of("JFK"), + List.of(List.of("JFK", List.of("BOS", "ORD")), List.of("BOS", List.of("JFK", "PWM"))))); } // ==================== Filter Tests ==================== @@ -433,10 +459,10 @@ public void testEmployeeHierarchyWithFilter() throws IOException { result, rows("Dev", "Eliot", 1, Collections.emptyList()), rows("Eliot", "Ron", 2, Collections.emptyList()), - rows("Ron", "Andrew", 3, List.of("{Andrew, null, 4}")), + rows("Ron", "Andrew", 3, List.of(Arrays.asList("Andrew", null, 4))), rows("Andrew", null, 4, Collections.emptyList()), rows("Asya", "Ron", 5, Collections.emptyList()), - rows("Dan", "Andrew", 6, List.of("{Andrew, null, 4}"))); + rows("Dan", "Andrew", 6, List.of(Arrays.asList("Andrew", null, 4)))); } /** @@ -496,7 +522,9 @@ public void testEmployeeHierarchyWithFilterAndMaxDepth() throws IOException { schema("reportingHierarchy", "array")); // Dev: reportsTo=Eliot -> Eliot(id=2) passes -> then Eliot.reportsTo=Ron -> Ron(id=3) passes // -> then Ron.reportsTo=Andrew -> Andrew(id=4) is filtered out -> stops - verifyDataRows(result, rows("Dev", "Eliot", 1, List.of("{Eliot, Ron, 2}", "{Ron, Andrew, 3}"))); + verifyDataRows( + result, + rows("Dev", "Eliot", 1, List.of(List.of("Eliot", "Ron", 2), List.of("Ron", "Andrew", 3)))); } // ==================== Edge Cases ==================== @@ -593,12 +621,12 @@ public void testGraphLookupWithFieldsProjection() throws IOException { verifySchema(result, schema("name", "string"), schema("reportingHierarchy", "array")); verifyDataRows( result, - rows("Dev", List.of("{Eliot, Ron, 2}")), - rows("Eliot", List.of("{Ron, Andrew, 3}")), - rows("Ron", List.of("{Andrew, null, 4}")), + rows("Dev", List.of(List.of("Eliot", "Ron", 2))), + rows("Eliot", List.of(List.of("Ron", "Andrew", 3))), + rows("Ron", List.of(Arrays.asList("Andrew", null, 4))), rows("Andrew", Collections.emptyList()), - rows("Asya", List.of("{Ron, Andrew, 3}")), - rows("Dan", List.of("{Andrew, null, 4}"))); + rows("Asya", List.of(List.of("Ron", "Andrew", 3))), + rows("Dan", List.of(Arrays.asList("Andrew", null, 4)))); } // ==================== Batch Mode Tests ==================== @@ -631,8 +659,8 @@ public void testBatchModeEmployeeHierarchy() throws IOException { verifyDataRows( result, rows( - List.of("{Dev, Eliot, 1}", "{Asya, Ron, 5}"), - List.of("{Ron, Andrew, 3, 0}", "{Andrew, null, 4, 1}"))); + List.of(List.of("Dev", "Eliot", 1), List.of("Asya", "Ron", 5)), + List.of(List.of("Ron", "Andrew", 3, 0), Arrays.asList("Andrew", null, 4, 1)))); } /** @@ -664,8 +692,11 @@ public void testBatchModeTravelersAirports() throws IOException { verifyDataRows( result, rows( - List.of("{Dev, JFK}", "{Eliot, JFK}", "{Jeff, BOS}"), - List.of("{JFK, [BOS, ORD], 0}", "{BOS, [JFK, PWM], 0}", "{PWM, [BOS, LHR], 1}"))); + List.of(List.of("Dev", "JFK"), List.of("Eliot", "JFK"), List.of("Jeff", "BOS")), + List.of( + List.of("JFK", List.of("BOS", "ORD"), 0), + List.of("BOS", List.of("JFK", "PWM"), 0), + List.of("PWM", List.of("BOS", "LHR"), 1)))); } /** @@ -696,12 +727,12 @@ public void testBatchModeBidirectional() throws IOException { verifyDataRows( result, rows( - List.of("{Dev, Eliot, 1}", "{Dan, Andrew, 6}"), + List.of(List.of("Dev", "Eliot", 1), List.of("Dan", "Andrew", 6)), List.of( - "{Dev, Eliot, 1, 0}", - "{Eliot, Ron, 2, 0}", - "{Andrew, null, 4, 0}", - "{Dan, Andrew, 6, 0}", - "{Asya, Ron, 5, 1}"))); + List.of("Dev", "Eliot", 1, 0), + List.of("Eliot", "Ron", 2, 0), + Arrays.asList("Andrew", null, 4, 0), + List.of("Dan", "Andrew", 6, 0), + List.of("Asya", "Ron", 5, 1)))); } } diff --git a/opensearch/src/main/java/org/opensearch/sql/opensearch/executor/OpenSearchExecutionEngine.java b/opensearch/src/main/java/org/opensearch/sql/opensearch/executor/OpenSearchExecutionEngine.java index 926fa4bc9b7..1f0f6d3fabf 100644 --- a/opensearch/src/main/java/org/opensearch/sql/opensearch/executor/OpenSearchExecutionEngine.java +++ b/opensearch/src/main/java/org/opensearch/sql/opensearch/executor/OpenSearchExecutionEngine.java @@ -11,6 +11,7 @@ import java.sql.ResultSetMetaData; import java.sql.SQLException; import java.util.ArrayList; +import java.util.Arrays; import java.util.HashMap; import java.util.LinkedHashMap; import java.util.List; @@ -247,7 +248,7 @@ private static Object processValue(Object value) throws SQLException { return convertedMap; } if (value instanceof StructImpl) { - return List.of(((StructImpl) value).getAttributes()); + return Arrays.asList(((StructImpl) value).getAttributes()); } if (value instanceof List) { List list = (List) value; From a7861c4a4c277c3f6633322691f59b348c775381 Mon Sep 17 00:00:00 2001 From: Heng Qian Date: Wed, 11 Feb 2026 18:25:57 +0800 Subject: [PATCH 4/4] Add experimental tag in doc Signed-off-by: Heng Qian --- docs/user/ppl/cmd/graphlookup.md | 2 +- docs/user/ppl/index.md | 4 +++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/docs/user/ppl/cmd/graphlookup.md b/docs/user/ppl/cmd/graphlookup.md index 0bb8ba3d196..2d6220edae4 100644 --- a/docs/user/ppl/cmd/graphlookup.md +++ b/docs/user/ppl/cmd/graphlookup.md @@ -1,5 +1,5 @@ -# graphLookup +# graphLookup (Experimental) The `graphLookup` command performs recursive graph traversal on a collection using a breadth-first search (BFS) algorithm. It searches for documents matching a start value and recursively traverses connections between documents based on specified fields. This is useful for hierarchical data like organizational charts, social networks, or routing graphs. diff --git a/docs/user/ppl/index.md b/docs/user/ppl/index.md index 12afe96eea0..718aa51f0fe 100644 --- a/docs/user/ppl/index.md +++ b/docs/user/ppl/index.md @@ -82,6 +82,8 @@ source=accounts | [addcoltotals command](cmd/addcoltotals.md) | 3.5 | stable (since 3.5) | Adds column values and appends a totals row. | | [transpose command](cmd/transpose.md) | 3.5 | stable (since 3.5) | Transpose rows to columns. | | [mvcombine command](cmd/mvcombine.md) | 3.5 | stable (since 3.4) | Combines values of a specified field across rows identical on all other fields. | +| [graphlookup command](cmd/graphlookup.md) | 3.5 | experimental (since 3.5) | Performs recursive graph traversal on a collection using a BFS algorithm.| + - [Syntax](cmd/syntax.md) - PPL query structure and command syntax formatting * **Functions** @@ -101,4 +103,4 @@ source=accounts * **Optimization** - [Optimization](../../user/optimization/optimization.rst) * **Limitations** - - [Limitations](limitations/limitations.md) \ No newline at end of file + - [Limitations](limitations/limitations.md)