From e1b33333de14c105ea653e41c62799e94647b8be Mon Sep 17 00:00:00 2001 From: Louis Chu Date: Mon, 20 Jul 2026 22:54:05 +0800 Subject: [PATCH] Add PPL multikv command (fixed-schema) Row-multiplying command that extracts fields from table-formatted text in an input field. Three-layer Calcite rewrite: MULTIKV_SPLIT UDF -> mvexpand (Uncollect/Correlate) -> per-column MULTIKV_EXTRACT UDF -> project. Output columns are resolved at plan time from a fields clause, forceheader, or positional noheader; a bare auto-header form is rejected with guidance. All columns emit VARCHAR (implicit per-op coercion downstream). Signed-off-by: Louis Chu --- .../org/opensearch/sql/analysis/Analyzer.java | 6 + .../sql/ast/AbstractNodeVisitor.java | 5 + .../org/opensearch/sql/ast/tree/Multikv.java | 89 ++++++++++ .../sql/calcite/CalciteRelNodeVisitor.java | 142 ++++++++++++++++ .../function/BuiltinFunctionName.java | 2 + .../function/PPLBuiltinOperators.java | 8 + .../expression/function/PPLFuncImpTable.java | 4 + .../multikv/MultikvExtractFunctionImpl.java | 68 ++++++++ .../function/multikv/MultikvParser.java | 136 +++++++++++++++ .../multikv/MultikvSplitFunctionImpl.java | 88 ++++++++++ docs/user/ppl/cmd/multikv.md | 95 +++++++++++ docs/user/ppl/index.md | 1 + .../remote/CalcitePPLMultikvCommandIT.java | 155 ++++++++++++++++++ .../sql/ppl/NewAddedCommandsIT.java | 29 ++++ ppl/src/main/antlr/OpenSearchPPLLexer.g4 | 4 + ppl/src/main/antlr/OpenSearchPPLParser.g4 | 17 ++ .../opensearch/sql/ppl/parser/AstBuilder.java | 33 ++++ .../sql/ppl/utils/PPLQueryDataAnonymizer.java | 20 +++ .../ppl/calcite/CalcitePPLMultikvTest.java | 94 +++++++++++ .../sql/ppl/parser/AstBuilderTest.java | 20 +++ .../ppl/utils/PPLQueryDataAnonymizerTest.java | 19 +++ 21 files changed, 1035 insertions(+) create mode 100644 core/src/main/java/org/opensearch/sql/ast/tree/Multikv.java create mode 100644 core/src/main/java/org/opensearch/sql/expression/function/multikv/MultikvExtractFunctionImpl.java create mode 100644 core/src/main/java/org/opensearch/sql/expression/function/multikv/MultikvParser.java create mode 100644 core/src/main/java/org/opensearch/sql/expression/function/multikv/MultikvSplitFunctionImpl.java create mode 100644 docs/user/ppl/cmd/multikv.md create mode 100644 integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePPLMultikvCommandIT.java create mode 100644 ppl/src/test/java/org/opensearch/sql/ppl/calcite/CalcitePPLMultikvTest.java diff --git a/core/src/main/java/org/opensearch/sql/analysis/Analyzer.java b/core/src/main/java/org/opensearch/sql/analysis/Analyzer.java index 701d1545b7..67ae51e468 100644 --- a/core/src/main/java/org/opensearch/sql/analysis/Analyzer.java +++ b/core/src/main/java/org/opensearch/sql/analysis/Analyzer.java @@ -83,6 +83,7 @@ import org.opensearch.sql.ast.tree.Lookup; import org.opensearch.sql.ast.tree.ML; import org.opensearch.sql.ast.tree.MakeResults; +import org.opensearch.sql.ast.tree.Multikv; import org.opensearch.sql.ast.tree.Multisearch; import org.opensearch.sql.ast.tree.MvCombine; import org.opensearch.sql.ast.tree.MvExpand; @@ -568,6 +569,11 @@ public LogicalPlan visitMakeResults(MakeResults node, AnalysisContext context) { throw getOnlyForCalciteException("makeresults"); } + @Override + public LogicalPlan visitMultikv(Multikv node, AnalysisContext context) { + throw getOnlyForCalciteException("multikv"); + } + @Override public LogicalPlan visitMvExpand(MvExpand node, AnalysisContext context) { throw getOnlyForCalciteException("mvexpand"); diff --git a/core/src/main/java/org/opensearch/sql/ast/AbstractNodeVisitor.java b/core/src/main/java/org/opensearch/sql/ast/AbstractNodeVisitor.java index acb6e10566..12591af1dd 100644 --- a/core/src/main/java/org/opensearch/sql/ast/AbstractNodeVisitor.java +++ b/core/src/main/java/org/opensearch/sql/ast/AbstractNodeVisitor.java @@ -72,6 +72,7 @@ import org.opensearch.sql.ast.tree.Lookup; import org.opensearch.sql.ast.tree.ML; import org.opensearch.sql.ast.tree.MakeResults; +import org.opensearch.sql.ast.tree.Multikv; import org.opensearch.sql.ast.tree.Multisearch; import org.opensearch.sql.ast.tree.MvCombine; import org.opensearch.sql.ast.tree.MvExpand; @@ -517,6 +518,10 @@ public T visitMvExpand(MvExpand node, C context) { return visitChildren(node, context); } + public T visitMultikv(Multikv node, C context) { + return visitChildren(node, context); + } + public T visitGraphLookup(GraphLookup node, C context) { return visitChildren(node, context); } diff --git a/core/src/main/java/org/opensearch/sql/ast/tree/Multikv.java b/core/src/main/java/org/opensearch/sql/ast/tree/Multikv.java new file mode 100644 index 0000000000..a8d7a76c68 --- /dev/null +++ b/core/src/main/java/org/opensearch/sql/ast/tree/Multikv.java @@ -0,0 +1,89 @@ +/* + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + */ + +package org.opensearch.sql.ast.tree; + +import com.google.common.collect.ImmutableList; +import java.util.List; +import javax.annotation.Nullable; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.ToString; +import org.opensearch.sql.ast.AbstractNodeVisitor; +import org.opensearch.sql.ast.expression.Field; + +/** + * AST node representing the {@code multikv} PPL command. + * + *

{@code multikv} extracts field values from an input field (default {@code _raw}) and emits one + * row per source row. The input field is either table-formatted text (split into columns) or an + * array of objects (one row per element, each declared column read from the element). This is a + * one-to-many (row-multiplying) command. + * + *

The output column names must be determinable at plan time, sourced from the declared {@link + * #fields} list, a literal {@link #forceHeader} line, or positional naming when {@link #noHeader} + * is set. Runtime header auto-detection (no fields, no forceheader, no noheader) is not supported; + * such a query is rejected at the field-resolution phase with a message directing the author to add + * a {@code fields} clause. + */ +@ToString +@EqualsAndHashCode(callSuper = false) +@Getter +public class Multikv extends UnresolvedPlan { + + /** Default Splunk input field for multikv. */ + public static final String DEFAULT_INPUT_FIELD = "_raw"; + + private UnresolvedPlan child; + + /** Input field carrying the table text. Defaults to {@code _raw}. */ + private final String inField; + + /** Declared output columns (the {@code fields} option). Null when not declared. */ + @Nullable private final List fields; + + /** Filter terms; a table row is kept only if it contains at least one term. Null when absent. */ + @Nullable private final List filterTerms; + + /** 1-based header line to force (the {@code forceheader} option). Null when absent. */ + @Nullable private final Integer forceHeader; + + /** When true, columns are named positionally (Column_1, Column_2, ...). */ + private final boolean noHeader; + + /** When true (default), the original event is dropped from the output. */ + private final boolean rmOrig; + + public Multikv( + String inField, + @Nullable List fields, + @Nullable List filterTerms, + @Nullable Integer forceHeader, + boolean noHeader, + boolean rmOrig) { + this.inField = inField; + this.fields = fields; + this.filterTerms = filterTerms; + this.forceHeader = forceHeader; + this.noHeader = noHeader; + this.rmOrig = rmOrig; + } + + @Override + public Multikv attach(UnresolvedPlan child) { + this.child = child; + return this; + } + + @Override + public List getChild() { + return this.child == null ? ImmutableList.of() : ImmutableList.of(this.child); + } + + @Override + public T accept(AbstractNodeVisitor nodeVisitor, C context) { + return nodeVisitor.visitMultikv(this, context); + } +} 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 94c40e5adb..a23c0af721 100644 --- a/core/src/main/java/org/opensearch/sql/calcite/CalciteRelNodeVisitor.java +++ b/core/src/main/java/org/opensearch/sql/calcite/CalciteRelNodeVisitor.java @@ -142,6 +142,7 @@ import org.opensearch.sql.ast.tree.Lookup.OutputStrategy; import org.opensearch.sql.ast.tree.ML; import org.opensearch.sql.ast.tree.MakeResults; +import org.opensearch.sql.ast.tree.Multikv; import org.opensearch.sql.ast.tree.Multisearch; import org.opensearch.sql.ast.tree.MvCombine; import org.opensearch.sql.ast.tree.MvExpand; @@ -197,6 +198,7 @@ import org.opensearch.sql.expression.function.BuiltinFunctionName; import org.opensearch.sql.expression.function.PPLBuiltinOperators; import org.opensearch.sql.expression.function.PPLFuncImpTable; +import org.opensearch.sql.expression.function.multikv.MultikvParser; import org.opensearch.sql.expression.parse.RegexCommonUtils; import org.opensearch.sql.utils.ParseUtils; import org.opensearch.sql.utils.WildcardRenameUtils; @@ -4469,6 +4471,146 @@ public RelNode visitMvExpand(MvExpand mvExpand, CalcitePlanContext context) { return relBuilder.peek(); } + /** + * multikv (fixed-schema): rewrite to an equivalent pipeline using existing operators. + * + *

+   *   eval __multikv_record__ = MULTIKV_SPLIT(inField, forceHeader, noHeader, filter)  // array<varchar>
+   *   | mvexpand __multikv_record__                                                     // 1 -> N rows
+   *   | eval <col> = MULTIKV_EXTRACT(__multikv_record__, '<col>')  for each declared field
+   *   | fields <col1>, <col2>, ...                                                      // declared output schema
+   * 
+ * + * The output column names come from the declared {@code fields} list and are therefore known at + * plan time. When no {@code fields} clause is declared, the output schema is not determinable at + * plan time and is rejected with guidance to add a {@code fields} clause. + * + *

When the input field is an object or an array of objects instead of text, the command + * dispatches to a native rewrite: {@code mvexpand} an array (a single object needs no explosion), + * then read each declared column from the object with {@code ITEM}. This preserves value types + * and shares the {@code fields} contract. + */ + @Override + public RelNode visitMultikv(Multikv node, CalcitePlanContext context) { + List fields = node.getFields(); + boolean noFields = (fields == null || fields.isEmpty()); + + // Fixed-schema: named output columns must be declared via the fields clause. The only + // no-fields form that is supported is positional noheader, which produces row-explosion with + // no named columns (e.g. `... | multikv noheader=true | stats count`). Any other no-fields + // form (bare auto-header, or forceheader without fields) cannot yield a plan-time schema and + // is rejected here with guidance, rather than failing at execution. + if (noFields && !node.isNoHeader()) { + throw ErrorReport.wrap( + new SemanticCheckException( + "multikv has no declared output columns. Add an explicit fields clause, for" + + " example: multikv fields ")) + .code(ErrorCode.FIELD_NOT_FOUND) + .location("while resolving the output schema for multikv") + .context("command", "multikv") + .build(); + } + + // Dispatch on the input field's type. A structured (array of objects) input is exploded + // natively with mvexpand and each declared column is read with ITEM, preserving element types. + // A text input falls through to the split pipeline below. The child is built once here to read + // its schema; the structured branch reuses that build, the text branch discards it. + boolean savedProjectVisited = context.isProjectVisited(); + RelNode probe = node.getChild().get(0).accept(this, context); + RelDataTypeField probeField = probe.getRowType().getField(node.getInField(), true, false); + boolean structuredArray = + probeField != null + && (SqlTypeUtil.isArray(probeField.getType()) + || SqlTypeUtil.isMultiset(probeField.getType())); + boolean structuredMap = probeField != null && SqlTypeUtil.isMap(probeField.getType()); + if (structuredArray || structuredMap) { + RelBuilder relBuilder = context.relBuilder; + if (structuredArray) { + // Array of objects: one row per element. A single object (map) needs no explosion. + buildExpandRelNode( + relBuilder.field(node.getInField()), + node.getInField(), + node.getInField(), + null, + context); + } + if (noFields) { + return relBuilder.peek(); + } + List projected = new ArrayList<>(); + List names = new ArrayList<>(); + for (Field f : fields) { + String col = f.getField().toString(); + RexNode item = + PPLFuncImpTable.INSTANCE.resolve( + context.rexBuilder, + BuiltinFunctionName.INTERNAL_ITEM, + relBuilder.field(node.getInField()), + context.rexBuilder.makeLiteral( + col, + context.rexBuilder.getTypeFactory().createSqlType(SqlTypeName.VARCHAR), + true)); + projected.add(item); + names.add(col); + } + relBuilder.project(projected, names); + context.setProjectVisited(true); + return relBuilder.peek(); + } + // Text input: discard the probe build and run the split pipeline on a fresh build. + context.relBuilder.build(); + context.setProjectVisited(savedProjectVisited); + + final String lineField = "__multikv_record__"; + final int forceHeader = node.getForceHeader() == null ? -1 : node.getForceHeader(); + final String filterJoined = + (node.getFilterTerms() == null || node.getFilterTerms().isEmpty()) + ? "" + : String.join(MultikvParser.FS, node.getFilterTerms()); + + // 1. eval __multikv_record__ = MULTIKV_SPLIT(inField, forceHeader, noHeader, filter) + UnresolvedPlan plan = + AstDSL.eval( + node.getChild().get(0), + AstDSL.let( + AstDSL.field(lineField), + AstDSL.function( + "multikv_split", + AstDSL.field(node.getInField()), + AstDSL.intLiteral(forceHeader), + AstDSL.booleanLiteral(node.isNoHeader()), + AstDSL.stringLiteral(filterJoined)))); + + // 2. mvexpand __multikv_record__ (one row per table data row) + plan = new MvExpand(AstDSL.field(lineField), null).attach(plan); + + if (noFields) { + // Positional noheader, no named columns: row-explosion only. The helper record column is + // retained (downstream typically only counts rows). Naming positional columns is deferred. + return plan.accept(this, context); + } + + // 3. eval each declared column = MULTIKV_EXTRACT(__multikv_record__, '') + Let[] lets = + fields.stream() + .map( + f -> { + String col = f.getField().toString(); + return AstDSL.let( + AstDSL.field(col), + AstDSL.function( + "multikv_extract", AstDSL.field(lineField), AstDSL.stringLiteral(col))); + }) + .toArray(Let[]::new); + plan = AstDSL.eval(plan, lets); + + // 4. project the declared columns only (drops original event fields and the helper record) + UnresolvedExpression[] projections = fields.toArray(new UnresolvedExpression[0]); + plan = AstDSL.project(plan, projections); + + return plan.accept(this, context); + } + @Override public RelNode visitValues(Values values, CalcitePlanContext context) { List> rows = values.getValues(); diff --git a/core/src/main/java/org/opensearch/sql/expression/function/BuiltinFunctionName.java b/core/src/main/java/org/opensearch/sql/expression/function/BuiltinFunctionName.java index e30d723ccf..ed023f49b0 100644 --- a/core/src/main/java/org/opensearch/sql/expression/function/BuiltinFunctionName.java +++ b/core/src/main/java/org/opensearch/sql/expression/function/BuiltinFunctionName.java @@ -272,6 +272,8 @@ public enum BuiltinFunctionName { JSON_ARRAY_LENGTH(FunctionName.of("json_array_length")), JSON_EXTRACT(FunctionName.of("json_extract")), JSON_EXTRACT_ALL(FunctionName.of("json_extract_all"), true), + MULTIKV_SPLIT(FunctionName.of("multikv_split"), true), + MULTIKV_EXTRACT(FunctionName.of("multikv_extract"), true), JSON_KEYS(FunctionName.of("json_keys")), JSON_SET(FunctionName.of("json_set")), JSON_DELETE(FunctionName.of("json_delete")), diff --git a/core/src/main/java/org/opensearch/sql/expression/function/PPLBuiltinOperators.java b/core/src/main/java/org/opensearch/sql/expression/function/PPLBuiltinOperators.java index d64f04bb9a..83400e5f52 100644 --- a/core/src/main/java/org/opensearch/sql/expression/function/PPLBuiltinOperators.java +++ b/core/src/main/java/org/opensearch/sql/expression/function/PPLBuiltinOperators.java @@ -141,6 +141,14 @@ public class PPLBuiltinOperators extends ReflectiveSqlOperatorTable { public static final SqlOperator JSON_APPEND = new JsonAppendFunctionImpl().toUDF("JSON_APPEND"); public static final SqlOperator JSON_EXTEND = new JsonExtendFunctionImpl().toUDF("JSON_EXTEND"); + // multikv internal functions + public static final SqlOperator MULTIKV_SPLIT = + new org.opensearch.sql.expression.function.multikv.MultikvSplitFunctionImpl() + .toUDF("MULTIKV_SPLIT"); + public static final SqlOperator MULTIKV_EXTRACT = + new org.opensearch.sql.expression.function.multikv.MultikvExtractFunctionImpl() + .toUDF("MULTIKV_EXTRACT"); + // Math functions public static final SqlOperator SPAN = new SpanFunction().toUDF("SPAN"); public static final SqlOperator E = new EulerFunction().toUDF("E"); diff --git a/core/src/main/java/org/opensearch/sql/expression/function/PPLFuncImpTable.java b/core/src/main/java/org/opensearch/sql/expression/function/PPLFuncImpTable.java index 151c4a9665..010bd0c2b6 100644 --- a/core/src/main/java/org/opensearch/sql/expression/function/PPLFuncImpTable.java +++ b/core/src/main/java/org/opensearch/sql/expression/function/PPLFuncImpTable.java @@ -163,6 +163,8 @@ import static org.opensearch.sql.expression.function.BuiltinFunctionName.MONTHNAME; import static org.opensearch.sql.expression.function.BuiltinFunctionName.MONTH_OF_YEAR; import static org.opensearch.sql.expression.function.BuiltinFunctionName.MSTIME; +import static org.opensearch.sql.expression.function.BuiltinFunctionName.MULTIKV_EXTRACT; +import static org.opensearch.sql.expression.function.BuiltinFunctionName.MULTIKV_SPLIT; import static org.opensearch.sql.expression.function.BuiltinFunctionName.MULTIMATCH; import static org.opensearch.sql.expression.function.BuiltinFunctionName.MULTIMATCHQUERY; import static org.opensearch.sql.expression.function.BuiltinFunctionName.MULTIPLY; @@ -1335,6 +1337,8 @@ void populate() { registerOperator(JSON_APPEND, PPLBuiltinOperators.JSON_APPEND); registerOperator(JSON_EXTEND, PPLBuiltinOperators.JSON_EXTEND); registerOperator(JSON_EXTRACT_ALL, PPLBuiltinOperators.JSON_EXTRACT_ALL); // internal + registerOperator(MULTIKV_SPLIT, PPLBuiltinOperators.MULTIKV_SPLIT); // internal + registerOperator(MULTIKV_EXTRACT, PPLBuiltinOperators.MULTIKV_EXTRACT); // internal // Register operators with a different type checker diff --git a/core/src/main/java/org/opensearch/sql/expression/function/multikv/MultikvExtractFunctionImpl.java b/core/src/main/java/org/opensearch/sql/expression/function/multikv/MultikvExtractFunctionImpl.java new file mode 100644 index 0000000000..f319e97fe2 --- /dev/null +++ b/core/src/main/java/org/opensearch/sql/expression/function/multikv/MultikvExtractFunctionImpl.java @@ -0,0 +1,68 @@ +/* + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + */ + +package org.opensearch.sql.expression.function.multikv; + +import static org.opensearch.sql.calcite.utils.OpenSearchTypeFactory.TYPE_FACTORY; + +import java.util.List; +import org.apache.calcite.adapter.enumerable.NotNullImplementor; +import org.apache.calcite.adapter.enumerable.NullPolicy; +import org.apache.calcite.adapter.enumerable.RexImpTable; +import org.apache.calcite.adapter.enumerable.RexToLixTranslator; +import org.apache.calcite.linq4j.tree.Expression; +import org.apache.calcite.linq4j.tree.Types; +import org.apache.calcite.rex.RexCall; +import org.apache.calcite.schema.impl.ScalarFunctionImpl; +import org.apache.calcite.sql.type.ReturnTypes; +import org.apache.calcite.sql.type.SqlReturnTypeInference; +import org.apache.calcite.sql.type.SqlTypeName; +import org.opensearch.sql.expression.function.ImplementorUDF; +import org.opensearch.sql.expression.function.UDFOperandMetadata; + +/** + * Internal UDF backing the {@code multikv} command. Extracts a single named cell value out of one + * serialized per-row record produced by {@link MultikvSplitFunctionImpl}. + * + *

Signature: {@code MULTIKV_EXTRACT(recordString, columnName)} returns {@code varchar} (null + * when the column is absent from the record). + */ +public class MultikvExtractFunctionImpl extends ImplementorUDF { + + public MultikvExtractFunctionImpl() { + super(new MultikvExtractImplementor(), NullPolicy.ANY); + } + + @Override + public SqlReturnTypeInference getReturnTypeInference() { + return ReturnTypes.explicit( + TYPE_FACTORY.createTypeWithNullability( + TYPE_FACTORY.createSqlType(SqlTypeName.VARCHAR), true)); + } + + @Override + public UDFOperandMetadata getOperandMetadata() { + return null; + } + + public static Object eval(Object... args) { + if (args.length < 2 || args[0] == null || args[1] == null) { + return null; + } + return MultikvParser.extract((String) args[0], (String) args[1]); + } + + public static class MultikvExtractImplementor implements NotNullImplementor { + @Override + public Expression implement( + RexToLixTranslator translator, RexCall call, List translatedOperands) { + ScalarFunctionImpl function = + (ScalarFunctionImpl) + ScalarFunctionImpl.create( + Types.lookupMethod(MultikvExtractFunctionImpl.class, "eval", Object[].class)); + return function.getImplementor().implement(translator, call, RexImpTable.NullAs.NULL); + } + } +} diff --git a/core/src/main/java/org/opensearch/sql/expression/function/multikv/MultikvParser.java b/core/src/main/java/org/opensearch/sql/expression/function/multikv/MultikvParser.java new file mode 100644 index 0000000000..329467d7db --- /dev/null +++ b/core/src/main/java/org/opensearch/sql/expression/function/multikv/MultikvParser.java @@ -0,0 +1,136 @@ +/* + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + */ + +package org.opensearch.sql.expression.function.multikv; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.regex.Pattern; + +/** + * Shared parsing logic for the {@code multikv} command UDFs. + * + *

{@link #parse} turns the table text of a single event into a list of serialized records, one + * per table data row. Each record encodes the row's cells as {@code namevalue} pairs joined by + * {@code }, where the column name comes from the detected header (or positional {@code + * Column_N} naming under {@code noheader}). {@link #extract} looks a single column value up out of + * one serialized record. + * + *

Columns are split on runs of whitespace. Richer aligned-offset handling is a later + * enhancement. + */ +public final class MultikvParser { + + /** Field separator between cells inside a serialized record. */ + public static final String FS = "\u001f"; + + /** Key/value separator between a column name and its value. */ + public static final String KV = "\u0002"; + + private static final Pattern WS = Pattern.compile("\\s+"); + private static final Pattern LINE = Pattern.compile("\\R"); + private static final Pattern FS_SPLIT = Pattern.compile(Pattern.quote(FS)); + + private MultikvParser() {} + + /** + * Parse the table text into serialized per-row records. + * + * @param raw the table text (typically the {@code _raw} field) + * @param forceHeader 1-based header line to force, or a value <= 0 for auto-detection + * @param noHeader when true, columns are named positionally ({@code Column_1}, ...) + * @param filterTerms keep only rows containing at least one term; empty means keep all + */ + public static List parse( + String raw, int forceHeader, boolean noHeader, List filterTerms) { + if (raw == null) { + return Collections.emptyList(); + } + List lines = new ArrayList<>(); + for (String l : LINE.split(raw)) { + if (!l.trim().isEmpty()) { + lines.add(l); + } + } + if (lines.isEmpty()) { + return Collections.emptyList(); + } + + List header; + int dataStart; + if (noHeader) { + header = null; + dataStart = 0; + } else if (forceHeader > 0) { + int hIdx = Math.min(forceHeader - 1, lines.size() - 1); + header = splitCols(lines.get(hIdx)); + dataStart = hIdx + 1; + } else { + header = splitCols(lines.get(0)); + dataStart = 1; + } + + boolean hasFilter = filterTerms != null && !filterTerms.isEmpty(); + List out = new ArrayList<>(); + for (int i = dataStart; i < lines.size(); i++) { + String line = lines.get(i); + if (hasFilter && !matchesFilter(line, filterTerms)) { + continue; + } + out.add(serialize(header, splitCols(line))); + } + return out; + } + + /** Extract the value for {@code col} out of one serialized record, or null when absent. */ + public static String extract(String record, String col) { + if (record == null || col == null) { + return null; + } + for (String pair : FS_SPLIT.split(record, -1)) { + int idx = pair.indexOf(KV); + if (idx < 0) { + continue; + } + if (pair.substring(0, idx).equals(col)) { + return pair.substring(idx + 1); + } + } + return null; + } + + private static boolean matchesFilter(String line, List filterTerms) { + for (String t : filterTerms) { + if (t != null && !t.isEmpty() && line.contains(t)) { + return true; + } + } + return false; + } + + private static List splitCols(String line) { + String t = line.trim(); + if (t.isEmpty()) { + return Collections.emptyList(); + } + return Arrays.asList(WS.split(t)); + } + + private static String serialize(List header, List cells) { + int n = (header != null) ? Math.max(header.size(), cells.size()) : cells.size(); + StringBuilder sb = new StringBuilder(); + for (int i = 0; i < n; i++) { + String name = (header != null && i < header.size()) ? header.get(i) : ("Column_" + (i + 1)); + String val = (i < cells.size()) ? cells.get(i) : ""; + if (sb.length() > 0) { + sb.append(FS); + } + sb.append(name).append(KV).append(val); + } + return sb.toString(); + } +} diff --git a/core/src/main/java/org/opensearch/sql/expression/function/multikv/MultikvSplitFunctionImpl.java b/core/src/main/java/org/opensearch/sql/expression/function/multikv/MultikvSplitFunctionImpl.java new file mode 100644 index 0000000000..f8d6736368 --- /dev/null +++ b/core/src/main/java/org/opensearch/sql/expression/function/multikv/MultikvSplitFunctionImpl.java @@ -0,0 +1,88 @@ +/* + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + */ + +package org.opensearch.sql.expression.function.multikv; + +import static org.apache.calcite.sql.type.SqlTypeUtil.createArrayType; + +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import org.apache.calcite.adapter.enumerable.NotNullImplementor; +import org.apache.calcite.adapter.enumerable.NullPolicy; +import org.apache.calcite.adapter.enumerable.RexImpTable; +import org.apache.calcite.adapter.enumerable.RexToLixTranslator; +import org.apache.calcite.linq4j.tree.Expression; +import org.apache.calcite.linq4j.tree.Types; +import org.apache.calcite.rel.type.RelDataTypeFactory; +import org.apache.calcite.rex.RexCall; +import org.apache.calcite.schema.impl.ScalarFunctionImpl; +import org.apache.calcite.sql.type.SqlReturnTypeInference; +import org.apache.calcite.sql.type.SqlTypeName; +import org.opensearch.sql.expression.function.ImplementorUDF; +import org.opensearch.sql.expression.function.UDFOperandMetadata; + +/** + * Internal UDF backing the {@code multikv} command. Parses the table text of one event into an + * array of serialized per-row records (see {@link MultikvParser}). The array is then exploded into + * one row per table data row (via the mvexpand uncollect/correlate primitive), after which {@link + * MultikvExtractFunctionImpl} pulls named cell values out of each record. + * + *

Signature: {@code MULTIKV_SPLIT(rawText, forceHeaderInt, noHeaderBool, filterJoined)} where + * {@code forceHeaderInt} is the 1-based forced header line or a value <= 0 for auto, and {@code + * filterJoined} is the filter terms joined by {@link MultikvParser#FS} (empty for none). Returns + * {@code array}. + */ +public class MultikvSplitFunctionImpl extends ImplementorUDF { + + public MultikvSplitFunctionImpl() { + super(new MultikvSplitImplementor(), NullPolicy.ANY); + } + + @Override + public SqlReturnTypeInference getReturnTypeInference() { + return sqlOperatorBinding -> { + RelDataTypeFactory typeFactory = sqlOperatorBinding.getTypeFactory(); + return createArrayType( + typeFactory, + typeFactory.createTypeWithNullability( + typeFactory.createSqlType(SqlTypeName.VARCHAR), true), + true); + }; + } + + @Override + public UDFOperandMetadata getOperandMetadata() { + return null; + } + + public static Object eval(Object... args) { + if (args.length < 1 || args[0] == null) { + return Collections.emptyList(); + } + String raw = (String) args[0]; + int forceHeader = (args.length > 1 && args[1] != null) ? ((Number) args[1]).intValue() : -1; + boolean noHeader = args.length > 2 && Boolean.TRUE.equals(args[2]); + String filterJoined = (args.length > 3 && args[3] != null) ? (String) args[3] : ""; + List filterTerms = + filterJoined.isEmpty() + ? Collections.emptyList() + : Arrays.asList( + filterJoined.split(java.util.regex.Pattern.quote(MultikvParser.FS), -1)); + return MultikvParser.parse(raw, forceHeader, noHeader, filterTerms); + } + + public static class MultikvSplitImplementor implements NotNullImplementor { + @Override + public Expression implement( + RexToLixTranslator translator, RexCall call, List translatedOperands) { + ScalarFunctionImpl function = + (ScalarFunctionImpl) + ScalarFunctionImpl.create( + Types.lookupMethod(MultikvSplitFunctionImpl.class, "eval", Object[].class)); + return function.getImplementor().implement(translator, call, RexImpTable.NullAs.NULL); + } + } +} diff --git a/docs/user/ppl/cmd/multikv.md b/docs/user/ppl/cmd/multikv.md new file mode 100644 index 0000000000..e04b82b860 --- /dev/null +++ b/docs/user/ppl/cmd/multikv.md @@ -0,0 +1,95 @@ + +# multikv + +The `multikv` command extracts field values from an input field and emits one row per source row. The input field can be table-formatted text (for example the aligned output of `ps`, `top`, `netstat`, or `df`), where a header row names the columns and each following line becomes its own output row. The input field can also be an object or an array of objects: a single object yields one row, an array yields one row per element, and each declared column is read from the object. Either way it is a one-to-many (row-multiplying) command. + +> **Note**: `multikv` is a streaming (mid-pipeline) command that runs on the coordinating node. It requires the Calcite engine (`plugins.calcite.enabled=true`). It reads from the input field named by `field=`, defaulting to `_raw`. When that field is text it is split into columns; when it is an array of objects each declared column is read from the element (`field= fields ...`). OpenSearch has no implicit `_raw` field, so for the text form either pass `field=` or place the text in `_raw` first with `eval _raw=`. Version 1 is fixed-schema: the output column names must be determinable at plan time via the `fields` clause. + +## Syntax + +The `multikv` command has the following syntax: + +```syntax +multikv [field=] [fields ...] [forceheader=] [noheader=] +``` + +## Parameters + +| Parameter | Required/Optional | Description | +| --- | --- | --- | +| `field=` | Optional | The input field that holds the table text. Defaults to `_raw`. Use this to read the text directly from a field such as `message` without a preceding `eval _raw=`. | +| `fields ...` | Optional | Declares the output columns to extract by name. This is the only form that yields named columns. Each extracted column is typed as `string`. | +| `forceheader` | Optional | The 1-based line number to use as the header, which skips banner lines above it. Must be a positive integer. Used together with `fields`. | +| `noheader` | Optional | When `true`, there is no header row; the command performs row explosion only and does not produce named columns (for example to count data lines). Default is `false`. | + +### Column typing + +Every extracted column is typed as `string`, because the values come from splitting text and there is no index mapping to infer a type from. The typed engine coerces these string columns per operation, so numeric use works without an explicit cast (for example `stats avg(pctIdle)` computes numerically, and `where pctIdle > 100` compares numerically). To pin a hard type, cast downstream, for example `eval n = cast(pctIdle as double)`. + +## Example 1: Extract a single column + +The following query reads the table text from the `raw` field with `field=` and extracts the `pctIdle` column: + +```ppl +source=metrics +| multikv field=raw fields pctIdle +| fields pctIdle +``` + +The query returns one row per table data row, with a single `pctIdle` (string) column. + +## Example 2: Extract multiple columns + +```ppl +source=metrics +| eval _raw = raw +| multikv fields CPU pctIdle +``` + +The query returns one row per table data row, with `CPU` (string) and `pctIdle` (string) columns. + +## Example 3: Skip a banner line with forceheader + +When the first line is a banner and the real header is on line 2: + +```ppl +source=report +| eval _raw = raw +| multikv fields endpoint forceheader=2 +``` + +The query uses line 2 as the header and returns the `endpoint` (string) column, one row per data line. + +## Example 4: Row explosion with noheader + +When there is no header row and only the number of data lines matters: + +```ppl +source=lines +| eval _raw = raw +| multikv noheader=true +| stats count +``` + +The query explodes the table text into one row per data line and counts them. + +## Example 5: Structured input (array of objects) + +When `field=` points at an array of objects, `multikv` emits one row per element and reads each declared column from the element, preserving the element value: + +```ppl +source=hosts +| multikv field=procs fields pid cpu +``` + +For a document with `procs = [{"pid":1,"cpu":0.5},{"pid":42,"cpu":9.1}]`, the query returns two rows: `(1, 0.5)` and `(42, 9.1)`. When `field=` points at a single object rather than an array, one row is returned. Nested container values are returned as-is; extract deeper fields downstream with `spath` or another `multikv field=`. + +## Limitations + +Version 1 is fixed-schema, so a bare `multikv` with no `fields` clause and no `noheader=true` cannot resolve its output column names at plan time and is rejected with guidance: + +```ppl +source=metrics | eval _raw = raw | multikv | fields pctIdle +``` + +Add an explicit `fields` clause to fix it, for example `multikv fields pctIdle`. Runtime header auto-detection is planned for a later version. The `filter` and `rmorig` options are not yet supported. diff --git a/docs/user/ppl/index.md b/docs/user/ppl/index.md index 1afb162963..8a6fadcefa 100644 --- a/docs/user/ppl/index.md +++ b/docs/user/ppl/index.md @@ -81,6 +81,7 @@ source=accounts | [explain command](cmd/explain.md) | 3.1 | stable (since 3.1) | Explain the plan of query. | | [show datasources command](cmd/showdatasources.md) | 2.4 | stable (since 2.4) | Query datasources configured in the PPL engine. | | [makeresults command](cmd/makeresults.md) | 3.8 | experimental (since 3.8) | Generate in-memory rows for testing and seeding, optionally from inline CSV/JSON data. | +| [multikv command](cmd/multikv.md) | 3.8 | experimental (since 3.8) | Extract fields from table-formatted text in a field, emitting one row per table data row. | | [addtotals command](cmd/addtotals.md) | 3.5 | stable (since 3.5) | Adds row and column values and appends a totals column and row. | | [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. | diff --git a/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePPLMultikvCommandIT.java b/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePPLMultikvCommandIT.java new file mode 100644 index 0000000000..f8e54c0c74 --- /dev/null +++ b/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePPLMultikvCommandIT.java @@ -0,0 +1,155 @@ +/* + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + */ + +package org.opensearch.sql.calcite.remote; + +import static org.opensearch.sql.util.MatcherUtils.rows; +import static org.opensearch.sql.util.MatcherUtils.schema; +import static org.opensearch.sql.util.MatcherUtils.verifyDataRows; +import static org.opensearch.sql.util.MatcherUtils.verifySchema; + +import java.io.IOException; +import org.json.JSONObject; +import org.junit.jupiter.api.Test; +import org.opensearch.client.Request; +import org.opensearch.client.ResponseException; +import org.opensearch.sql.legacy.TestUtils; +import org.opensearch.sql.ppl.PPLIntegTestCase; + +/** Integration tests for the {@code multikv} command (fixed-schema). */ +public class CalcitePPLMultikvCommandIT extends PPLIntegTestCase { + + @Override + public void init() throws Exception { + super.init(); + enableCalcite(); + + // Auto-header table: header on line 1, two data rows. + if (!TestUtils.isIndexExist(client(), "test_multikv")) { + TestUtils.createIndexByRestClient(client(), "test_multikv", null); + Request doc = new Request("PUT", "/test_multikv/_doc/1?refresh=true"); + // raw = "CPU pctUser pctIdle\nall 5 90\n0 3 92" + doc.setJsonEntity("{\"raw\": \"CPU pctUser pctIdle\\nall 5 90\\n0 3 92\"}"); + client().performRequest(doc); + } + + // forceheader: line 1 is a junk banner, the real header is line 2. + if (!TestUtils.isIndexExist(client(), "test_multikv_force")) { + TestUtils.createIndexByRestClient(client(), "test_multikv_force", null); + Request doc = new Request("PUT", "/test_multikv_force/_doc/1?refresh=true"); + // raw = "== report ==\nendpoint count\nfoo 1\nbar 2" + doc.setJsonEntity("{\"raw\": \"== report ==\\nendpoint count\\nfoo 1\\nbar 2\"}"); + client().performRequest(doc); + } + + // noheader: 3 data lines, no header row. + if (!TestUtils.isIndexExist(client(), "test_multikv_noheader")) { + TestUtils.createIndexByRestClient(client(), "test_multikv_noheader", null); + Request doc = new Request("PUT", "/test_multikv_noheader/_doc/1?refresh=true"); + doc.setJsonEntity("{\"raw\": \"a 1\\nb 2\\nc 3\"}"); + client().performRequest(doc); + } + + // structured: procs is an array of objects (nested), two elements. + if (!TestUtils.isIndexExist(client(), "test_multikv_struct")) { + String mapping = + "{\"mappings\":{\"properties\":{\"procs\":{\"type\":\"nested\"," + + "\"properties\":{\"pid\":{\"type\":\"long\"},\"cpu\":{\"type\":\"double\"}}}}}}"; + TestUtils.createIndexByRestClient(client(), "test_multikv_struct", mapping); + Request doc = new Request("PUT", "/test_multikv_struct/_doc/1?refresh=true"); + doc.setJsonEntity("{\"procs\":[{\"pid\":1,\"cpu\":0.5},{\"pid\":42,\"cpu\":9.1}]}"); + client().performRequest(doc); + } + + // single object: meta is one object (map), not an array. + if (!TestUtils.isIndexExist(client(), "test_multikv_obj")) { + String mapping = + "{\"mappings\":{\"properties\":{\"meta\":{\"properties\":" + + "{\"owner\":{\"type\":\"keyword\"},\"role\":{\"type\":\"keyword\"}}}}}}"; + TestUtils.createIndexByRestClient(client(), "test_multikv_obj", mapping); + Request doc = new Request("PUT", "/test_multikv_obj/_doc/1?refresh=true"); + doc.setJsonEntity("{\"meta\":{\"owner\":\"root\",\"role\":\"admin\"}}"); + client().performRequest(doc); + } + } + + @Test + public void testMultikvFields() throws IOException { + // Declared single column: auto-detected header maps pctIdle -> its column. + JSONObject result = + executeQuery( + "source=test_multikv | eval _raw = raw | multikv fields pctIdle | fields pctIdle"); + verifySchema(result, schema("pctIdle", "string")); + verifyDataRows(result, rows("90"), rows("92")); + } + + @Test + public void testMultikvFieldsMultipleColumns() throws IOException { + JSONObject result = + executeQuery("source=test_multikv | eval _raw = raw | multikv fields CPU pctIdle"); + verifySchema(result, schema("CPU", "string"), schema("pctIdle", "string")); + verifyDataRows(result, rows("all", "90"), rows("0", "92")); + } + + @Test + public void testMultikvForceHeader() throws IOException { + // forceheader=2 skips the banner and uses line 2 ("endpoint count") as the header. + JSONObject result = + executeQuery( + "source=test_multikv_force | eval _raw = raw | multikv fields endpoint forceheader=2 |" + + " fields endpoint"); + verifySchema(result, schema("endpoint", "string")); + verifyDataRows(result, rows("foo"), rows("bar")); + } + + @Test + public void testMultikvNoHeaderCount() throws IOException { + // Positional noheader, downstream only counts rows (3 data lines -> count 3). + JSONObject result = + executeQuery( + "source=test_multikv_noheader | eval _raw = raw | multikv noheader=true | stats count"); + verifyDataRows(result, rows(3)); + } + + @Test + public void testMultikvFieldEquals() throws IOException { + // field= reads the table text directly from a field, no eval _raw needed. + JSONObject result = + executeQuery("source=test_multikv | multikv field=raw fields pctIdle | fields pctIdle"); + verifySchema(result, schema("pctIdle", "string")); + verifyDataRows(result, rows("90"), rows("92")); + } + + @Test + public void testMultikvStructuredArrayOfObjects() throws IOException { + // Structured input: procs is an array of objects. multikv explodes it into one row per + // element and reads each declared column with ITEM. + JSONObject result = + executeQuery( + "source=test_multikv_struct | multikv field=procs fields pid cpu" + + " | eval p = cast(pid as int) | fields p | sort p"); + verifyDataRows(result, rows(1), rows(42)); + } + + @Test + public void testMultikvSingleObject() throws IOException { + // Single object (map): no row multiplication, each declared column read with ITEM -> 1 row. + JSONObject result = + executeQuery("source=test_multikv_obj | multikv field=meta fields owner role"); + verifyDataRows(result, rows("root", "admin")); + } + + @Test + public void testBareMultikvRejectedWithGuidance() throws IOException { + // Bare auto-header multikv referencing a column downstream must be rejected at plan time + // with guidance to declare a fields clause. + Throwable t = + org.junit.Assert.assertThrows( + ResponseException.class, + () -> executeQuery("source=test_multikv | eval _raw = raw | multikv | fields pctIdle")); + org.junit.Assert.assertTrue( + t.getMessage(), t.getMessage().toLowerCase().contains("fields clause")); + } +} diff --git a/integ-test/src/test/java/org/opensearch/sql/ppl/NewAddedCommandsIT.java b/integ-test/src/test/java/org/opensearch/sql/ppl/NewAddedCommandsIT.java index 6b5ac0d430..127adc4230 100644 --- a/integ-test/src/test/java/org/opensearch/sql/ppl/NewAddedCommandsIT.java +++ b/integ-test/src/test/java/org/opensearch/sql/ppl/NewAddedCommandsIT.java @@ -18,6 +18,7 @@ import org.json.JSONArray; import org.json.JSONObject; import org.junit.jupiter.api.Test; +import org.opensearch.client.Request; import org.opensearch.client.ResponseException; import org.opensearch.sql.util.TestUtils; @@ -339,6 +340,34 @@ public void testMakeResults() throws IOException { } } + @Test + public void testMultikv() throws IOException { + // v3-marker (registered on Calcite, rejected on V2). The table text (with newlines) lives in + // indexed data, not the query string, so the request payload stays valid JSON: a header line + // plus two data rows yields two extracted rows. + String index = "test_new_added_multikv"; + if (!TestUtils.isIndexExist(client(), index)) { + TestUtils.createIndexByRestClient(client(), index, null); + Request doc = new Request("PUT", "/" + index + "/_doc/1?refresh=true"); + doc.setJsonEntity("{\"raw\": \"CPU pctIdle\\nall 90\\n0 92\"}"); + client().performRequest(doc); + } + JSONObject result; + try { + result = + executeQuery("source=" + index + " | multikv field=raw fields pctIdle | fields pctIdle"); + } catch (ResponseException e) { + result = new JSONObject(TestUtils.getResponseBody(e.getResponse())); + } + + if (isCalciteEnabled()) { + assertThat(result.getJSONArray("datarows").length(), equalTo(2)); + } else { + // multikv is Calcite-only, so V2 returns an error. + assertThat(result.has("error"), equalTo(true)); + } + } + @Test public void testMvExpandCommandBasicExpansion() throws IOException { JSONObject result; diff --git a/ppl/src/main/antlr/OpenSearchPPLLexer.g4 b/ppl/src/main/antlr/OpenSearchPPLLexer.g4 index b26751ad61..573e7177ad 100644 --- a/ppl/src/main/antlr/OpenSearchPPLLexer.g4 +++ b/ppl/src/main/antlr/OpenSearchPPLLexer.g4 @@ -76,6 +76,10 @@ ROW: 'ROW'; COL: 'COL'; EXPAND: 'EXPAND'; MVEXPAND: 'MVEXPAND'; +MULTIKV: 'MULTIKV'; +FORCEHEADER: 'FORCEHEADER'; +NOHEADER: 'NOHEADER'; +RMORIG: 'RMORIG'; SIMPLE_PATTERN: 'SIMPLE_PATTERN'; BRAIN: 'BRAIN'; VARIABLE_COUNT_THRESHOLD: 'VARIABLE_COUNT_THRESHOLD'; diff --git a/ppl/src/main/antlr/OpenSearchPPLParser.g4 b/ppl/src/main/antlr/OpenSearchPPLParser.g4 index eeaed6daf5..37895f4cfb 100644 --- a/ppl/src/main/antlr/OpenSearchPPLParser.g4 +++ b/ppl/src/main/antlr/OpenSearchPPLParser.g4 @@ -86,6 +86,7 @@ commands | appendCommand | expandCommand | mvexpandCommand + | multikvCommand | flattenCommand | reverseCommand | regexCommand @@ -135,6 +136,7 @@ commandName | CONVERT | EXPAND | MVEXPAND + | MULTIKV | FLATTEN | TRENDLINE | TIMECHART @@ -652,6 +654,21 @@ mvexpandCommand : MVEXPAND fieldExpression (LIMIT EQUAL INTEGER_LITERAL)? ; +// multikv: extract field values from table-formatted text in a field (default _raw), +// emitting one row per table data row. v1 fixed-schema: output columns must be +// determinable at plan time via the fields list, a literal forceheader, or noheader. +multikvCommand + : MULTIKV multikvParameter* + ; + +multikvParameter + : FIELDS fields = fieldList + | FIELD EQUAL inField = qualifiedName + | FORCEHEADER EQUAL forceHeader = integerLiteral + | NOHEADER EQUAL noHeader = booleanLiteral + | RMORIG EQUAL rmOrig = booleanLiteral + ; + flattenCommand : FLATTEN fieldExpression (AS aliases = identifierSeq)? ; 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 e87264909c..efe35f8baa 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 @@ -99,6 +99,7 @@ import org.opensearch.sql.ast.tree.ML; import org.opensearch.sql.ast.tree.MakeResults; import org.opensearch.sql.ast.tree.MinSpanBin; +import org.opensearch.sql.ast.tree.Multikv; import org.opensearch.sql.ast.tree.Multisearch; import org.opensearch.sql.ast.tree.MvCombine; import org.opensearch.sql.ast.tree.MvExpand; @@ -1076,6 +1077,38 @@ public UnresolvedPlan visitMvexpandCommand(OpenSearchPPLParser.MvexpandCommandCo return new MvExpand(field, limit); } + @Override + public UnresolvedPlan visitMultikvCommand(OpenSearchPPLParser.MultikvCommandContext ctx) { + List fields = null; + String inField = Multikv.DEFAULT_INPUT_FIELD; + Integer forceHeader = null; + boolean noHeader = false; + boolean rmOrig = true; // Splunk default + + for (OpenSearchPPLParser.MultikvParameterContext p : ctx.multikvParameter()) { + if (p.fields != null) { + fields = + p.fields.fieldExpression().stream() + .map(f -> (Field) expressionBuilder.visit(f)) + .collect(Collectors.toList()); + } else if (p.inField != null) { + inField = p.inField.getText(); + } else if (p.forceHeader != null) { + forceHeader = Integer.parseInt(p.forceHeader.getText()); + if (forceHeader <= 0) { + throw new IllegalArgumentException( + "multikv forceheader must be a positive line number, got: " + forceHeader); + } + } else if (p.noHeader != null) { + noHeader = Boolean.parseBoolean(p.noHeader.getText()); + } else if (p.rmOrig != null) { + rmOrig = Boolean.parseBoolean(p.rmOrig.getText()); + } + } + + return new Multikv(inField, fields, null, forceHeader, noHeader, rmOrig); + } + @Override public UnresolvedPlan visitGrokCommand(OpenSearchPPLParser.GrokCommandContext ctx) { UnresolvedExpression sourceField = internalVisitExpression(ctx.source_field); 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 11c47d137e..4629197886 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 @@ -84,6 +84,7 @@ import org.opensearch.sql.ast.tree.Join; import org.opensearch.sql.ast.tree.Lookup; import org.opensearch.sql.ast.tree.MinSpanBin; +import org.opensearch.sql.ast.tree.Multikv; import org.opensearch.sql.ast.tree.Multisearch; import org.opensearch.sql.ast.tree.MvCombine; import org.opensearch.sql.ast.tree.MvExpand; @@ -611,6 +612,25 @@ public String visitForeach(Foreach node, String context) { return command.append(" [ eval ").append(evalClauses).append(" ]").toString(); } + @Override + public String visitMultikv(Multikv node, String context) { + String child = node.getChild().get(0).accept(this, context); + StringBuilder sb = new StringBuilder(child).append(" | multikv"); + if (node.getFields() != null && !node.getFields().isEmpty()) { + sb.append(" fields ").append(MASK_COLUMN); + } + if (node.getFilterTerms() != null && !node.getFilterTerms().isEmpty()) { + sb.append(" filter ").append(MASK_LITERAL); + } + if (node.getForceHeader() != null) { + sb.append(" forceheader=").append(MASK_LITERAL); + } + if (node.isNoHeader()) { + sb.append(" noheader=true"); + } + return sb.toString(); + } + /** Build {@link LogicalSort}. */ @Override public String visitSort(Sort node, String context) { diff --git a/ppl/src/test/java/org/opensearch/sql/ppl/calcite/CalcitePPLMultikvTest.java b/ppl/src/test/java/org/opensearch/sql/ppl/calcite/CalcitePPLMultikvTest.java new file mode 100644 index 0000000000..abe55d6382 --- /dev/null +++ b/ppl/src/test/java/org/opensearch/sql/ppl/calcite/CalcitePPLMultikvTest.java @@ -0,0 +1,94 @@ +/* + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + */ + +package org.opensearch.sql.ppl.calcite; + +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; + +import org.apache.calcite.plan.RelOptUtil; +import org.apache.calcite.rel.RelNode; +import org.apache.calcite.test.CalciteAssert; +import org.junit.Test; + +/** + * Logical-plan tests for the {@code multikv} command (fixed-schema). + * + *

multikv lowers to a three-layer rewrite: an {@code eval} that calls {@code MULTIKV_SPLIT} into + * a helper record column, an mvexpand (Uncollect over a Correlate) that explodes one row per table + * data row, and a per-declared-column {@code MULTIKV_EXTRACT} eval followed by a narrowing project. + * These assertions pin the presence and absence of those operators rather than the full Rex + * rendering, which is exercised end-to-end by CalcitePPLMultikvCommandIT. + */ +public class CalcitePPLMultikvTest extends CalcitePPLAbstractTest { + public CalcitePPLMultikvTest() { + super(CalciteAssert.SchemaSpec.SCOTT_WITH_TEMPORAL); + } + + private String plan(String ppl) { + RelNode root = getRelNode(ppl); + return RelOptUtil.toString(root); + } + + private void expectError(String ppl, String messageFragment) { + try { + getRelNode(ppl); + fail("expected an error for: " + ppl); + } catch (Exception e) { + String msg = String.valueOf(e.getMessage()); + assertTrue( + "expected message containing '" + messageFragment + "' but got: " + msg, + msg.contains(messageFragment)); + } + } + + @Test + public void testMultikvFieldsSingleColumn() { + String p = plan("source=EMP | eval _raw = ENAME | multikv fields pctIdle | fields pctIdle"); + // parse layer + assertTrue(p, p.contains("MULTIKV_SPLIT")); + // explode layer (mvexpand = Uncollect over Correlate) + assertTrue(p, p.contains("Uncollect") || p.contains("Correlate")); + // extract layer for the declared column + assertTrue(p, p.contains("MULTIKV_EXTRACT")); + assertTrue(p, p.contains("pctIdle")); + } + + @Test + public void testMultikvFieldsMultipleColumns() { + String p = plan("source=EMP | eval _raw = ENAME | multikv fields CPU pctIdle"); + // one extract per declared column + int idx = p.indexOf("MULTIKV_EXTRACT"); + assertTrue(p, idx >= 0 && p.indexOf("MULTIKV_EXTRACT", idx + 1) >= 0); + assertTrue(p, p.contains("CPU")); + assertTrue(p, p.contains("pctIdle")); + } + + @Test + public void testMultikvForceHeader() { + // forceheader=2 is threaded into MULTIKV_SPLIT as the header-line argument. + String p = plan("source=EMP | eval _raw = ENAME | multikv fields endpoint forceheader=2"); + assertTrue(p, p.contains("MULTIKV_SPLIT")); + assertTrue(p, p.contains("MULTIKV_EXTRACT")); + assertTrue(p, p.contains("endpoint")); + } + + @Test + public void testMultikvNoHeaderIsRowExplosionOnly() { + // Positional noheader with no fields: split + explode, but no named-column extract/project. + String p = plan("source=EMP | eval _raw = ENAME | multikv noheader=true"); + assertTrue(p, p.contains("MULTIKV_SPLIT")); + assertTrue(p, p.contains("Uncollect") || p.contains("Correlate")); + assertFalse(p, p.contains("MULTIKV_EXTRACT")); + } + + @Test + public void testBareMultikvRejectedWithGuidance() { + // Bare auto-header multikv (no fields, no noheader) cannot resolve a plan-time schema and is + // rejected at the field-resolution phase with actionable guidance. + expectError("source=EMP | eval _raw = ENAME | multikv", "fields clause"); + } +} diff --git a/ppl/src/test/java/org/opensearch/sql/ppl/parser/AstBuilderTest.java b/ppl/src/test/java/org/opensearch/sql/ppl/parser/AstBuilderTest.java index d5f45a96f8..07afb12fb9 100644 --- a/ppl/src/test/java/org/opensearch/sql/ppl/parser/AstBuilderTest.java +++ b/ppl/src/test/java/org/opensearch/sql/ppl/parser/AstBuilderTest.java @@ -82,6 +82,7 @@ import org.opensearch.sql.ast.tree.Kmeans; import org.opensearch.sql.ast.tree.ML; import org.opensearch.sql.ast.tree.MakeResults; +import org.opensearch.sql.ast.tree.Multikv; import org.opensearch.sql.ast.tree.RareTopN.CommandType; import org.opensearch.sql.common.antlr.SyntaxCheckException; import org.opensearch.sql.common.setting.Settings.Key; @@ -1114,6 +1115,25 @@ public void testMakeResultsCommand() { assertEqual("makeresults count=5", new MakeResults(5)); } + @Test + public void testMultikvCommand() { + assertEqual( + "source=t | multikv fields CPU pctIdle", + new Multikv("_raw", Arrays.asList(field("CPU"), field("pctIdle")), null, null, false, true) + .attach(relation("t"))); + assertEqual( + "source=t | multikv fields endpoint forceheader=2", + new Multikv("_raw", Arrays.asList(field("endpoint")), null, 2, false, true) + .attach(relation("t"))); + assertEqual( + "source=t | multikv noheader=true", + new Multikv("_raw", null, null, null, true, true).attach(relation("t"))); + assertEqual( + "source=t | multikv field=message fields CPU", + new Multikv("message", Arrays.asList(field("CPU")), null, null, false, true) + .attach(relation("t"))); + } + @Test public void testDescribeMatchAllCrossClusterSearchCommand() { assertEqual("describe *:t", describe(mappingTable("*:t"))); 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 eba4f57112..63cc8f6c87 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 @@ -41,6 +41,25 @@ public void testMakeResultsCommand() { assertEquals("makeresults", anonymize("makeresults count=5")); } + @Test + public void testMultikvFieldsCommand() { + assertEquals( + "source=table | multikv fields identifier", anonymize("source=t | multikv fields pctIdle")); + } + + @Test + public void testMultikvForceHeaderCommand() { + assertEquals( + "source=table | multikv fields identifier forceheader=***", + anonymize("source=t | multikv fields endpoint forceheader=2")); + } + + @Test + public void testMultikvNoHeaderCommand() { + assertEquals( + "source=table | multikv noheader=true", anonymize("source=t | multikv noheader=true")); + } + @Test public void testTableFunctionCommand() { assertEquals(