From e146e91d15bebbfcc26bce1ee3099b198a6f3ffb Mon Sep 17 00:00:00 2001 From: Louis Chu Date: Mon, 20 Jul 2026 22:54:05 +0800 Subject: [PATCH 1/2] 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 | 139 ++++++++++++++++ .../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 | 97 +++++++++++ docs/user/ppl/index.md | 1 + .../remote/CalcitePPLMultikvCommandIT.java | 153 ++++++++++++++++++ .../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, 1032 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 701d1545b76..67ae51e4688 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 acb6e105661..12591af1dd3 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 00000000000..a8d7a76c68b --- /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 94c40e5adb2..0aa56c9b160 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,143 @@ 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}. The extracted columns are + * typed {@code ANY}: object and nested fields collapse to ANY-valued containers in the type + * layer, so the mapped scalar type is not recovered (cast downstream). 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: output columns must come from the fields clause. The only supported no-fields + // form is positional noheader (row-explosion, no named columns). Any other no-fields form + // (bare auto-header, or forceheader without fields) has no plan-time schema and is rejected. + 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; the extracted column is + // typed ANY (element types are erased to ANY upstream), not the mapped scalar type. 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()); + + 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)))); + + 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); + } + + 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); + + 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 e30d723ccfc..ed023f49b0e 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 d64f04bb9ad..83400e5f52e 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 151c4a96655..010bd0c2b60 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 00000000000..f319e97fe21 --- /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 00000000000..329467d7db3 --- /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 00000000000..f8d67363683 --- /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 00000000000..a42f1a0c9b3 --- /dev/null +++ b/docs/user/ppl/cmd/multikv.md @@ -0,0 +1,97 @@ + +# 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)`. + +For structured input (an object or an array of objects), the extracted columns are typed `ANY` rather than the field's mapped scalar type, because object and nested fields collapse to `ANY`-valued containers before the command runs, so the mapped type is not recovered. Cast downstream for typed operations, for example `eval p = cast(pid as int)`. + +## 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 1afb162963d..8a6fadcefa5 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 00000000000..7e06d15e1f2 --- /dev/null +++ b/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalcitePPLMultikvCommandIT.java @@ -0,0 +1,153 @@ +/* + * 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"); + 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"); + 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 6b5ac0d4302..127adc42307 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 b26751ad61b..573e7177ad5 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 eeaed6daf52..37895f4cfb1 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 e87264909c8..efe35f8baa9 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 11c47d137e8..46291978864 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 00000000000..abe55d63822 --- /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 d5f45a96f8c..07afb12fb93 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 eba4f57112b..63cc8f6c87c 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( From 3a0aee8837958946b75887926b3b9b1bef52184c Mon Sep 17 00:00:00 2001 From: Louis Chu Date: Tue, 11 Aug 2026 16:30:06 +0800 Subject: [PATCH 2/2] multikv: support inline fields col:type for plan-time typed columns Add an optional col:type declaration to the multikv fields clause so an extracted column is typed at plan time instead of the default (string in text mode, ANY in structured mode). The type is lowered as a plan-time safe cast, equivalent to hoisting | eval col = cast(col as ) into the command, and works in both the text and structured branches. - Grammar: new multikv-local multikvField rule; the typed form is matched via the CLUSTER token (the case-insensitive lexer folds col: into it) and the column name is that token minus its trailing colon. The shared fieldList is untouched. - Extract the makeresults type-name resolver into a shared PplInlineTypeResolver (ppl.utils) so both commands share one scalar vocabulary and one UDT-rejection policy (date/time/timestamp/ip/json -> 'use string and cast'). - Multikv AST node carries a parallel per-column fieldTypes list; visitMultikv wraps INTERNAL_ITEM (structured) and MULTIKV_EXTRACT (text) in a safe cast when a type is declared. - Tests: AstBuilder typed-fields parse, plan-shape (SAFE_CAST present, UDT rejected), and IT for typed text + typed structured columns. Note: a typed column name must start with a letter or * because col: lexes as the cross-cluster prefix token; otherwise declare it untyped and cast downstream. Signed-off-by: Louis Chu --- .../org/opensearch/sql/ast/tree/Multikv.java | 21 ++- .../sql/calcite/CalciteRelNodeVisitor.java | 121 ++++++++++++------ docs/category.json | 3 +- docs/user/ppl/cmd/multikv.md | 75 ++++++++--- docs/user/ppl/index.md | 2 +- doctest/test_data/multikv_lines.json | 1 + doctest/test_data/multikv_report.json | 1 + doctest/test_data/multikv_struct.json | 1 + doctest/test_data/multikv_text.json | 1 + doctest/test_docs.py | 4 + doctest/test_mapping/multikv_struct.json | 1 + .../sql/calcite/remote/CalciteExplainIT.java | 41 ++++++ .../remote/CalcitePPLMultikvCommandIT.java | 53 ++++++++ ppl/src/main/antlr/OpenSearchPPLParser.g4 | 14 +- .../opensearch/sql/ppl/parser/AstBuilder.java | 32 ++++- .../sql/ppl/utils/MakeResultsDataParser.java | 25 +--- .../sql/ppl/utils/PplInlineTypeResolver.java | 53 ++++++++ .../ppl/calcite/CalcitePPLMultikvTest.java | 25 ++++ .../sql/ppl/parser/AstBuilderTest.java | 28 ++++ 19 files changed, 412 insertions(+), 90 deletions(-) create mode 100644 doctest/test_data/multikv_lines.json create mode 100644 doctest/test_data/multikv_report.json create mode 100644 doctest/test_data/multikv_struct.json create mode 100644 doctest/test_data/multikv_text.json create mode 100644 doctest/test_mapping/multikv_struct.json create mode 100644 ppl/src/main/java/org/opensearch/sql/ppl/utils/PplInlineTypeResolver.java 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 index a8d7a76c68b..a16c0c16b0c 100644 --- a/core/src/main/java/org/opensearch/sql/ast/tree/Multikv.java +++ b/core/src/main/java/org/opensearch/sql/ast/tree/Multikv.java @@ -13,6 +13,7 @@ import lombok.ToString; import org.opensearch.sql.ast.AbstractNodeVisitor; import org.opensearch.sql.ast.expression.Field; +import org.opensearch.sql.data.type.ExprCoreType; /** * AST node representing the {@code multikv} PPL command. @@ -33,7 +34,7 @@ @Getter public class Multikv extends UnresolvedPlan { - /** Default Splunk input field for multikv. */ + /** Default input field name used when {@code field=} is omitted. */ public static final String DEFAULT_INPUT_FIELD = "_raw"; private UnresolvedPlan child; @@ -44,6 +45,12 @@ public class Multikv extends UnresolvedPlan { /** Declared output columns (the {@code fields} option). Null when not declared. */ @Nullable private final List fields; + /** + * Per-column declared types aligned with {@link #fields} (the {@code col:type} syntax). Null when + * none declared. + */ + @Nullable private final List fieldTypes; + /** Filter terms; a table row is kept only if it contains at least one term. Null when absent. */ @Nullable private final List filterTerms; @@ -63,8 +70,20 @@ public Multikv( @Nullable Integer forceHeader, boolean noHeader, boolean rmOrig) { + this(inField, fields, null, filterTerms, forceHeader, noHeader, rmOrig); + } + + public Multikv( + String inField, + @Nullable List fields, + @Nullable List fieldTypes, + @Nullable List filterTerms, + @Nullable Integer forceHeader, + boolean noHeader, + boolean rmOrig) { this.inField = inField; this.fields = fields; + this.fieldTypes = fieldTypes; this.filterTerms = filterTerms; this.forceHeader = forceHeader; this.noHeader = noHeader; 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 2ba89889a61..88135c973b4 100644 --- a/core/src/main/java/org/opensearch/sql/calcite/CalciteRelNodeVisitor.java +++ b/core/src/main/java/org/opensearch/sql/calcite/CalciteRelNodeVisitor.java @@ -4705,12 +4705,10 @@ public RelNode visitMultikv(Multikv node, CalcitePlanContext context) { .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; the extracted column is - // typed ANY (element types are erased to ANY upstream), not the mapped scalar type. 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(); + // Dispatch on the input field's type. The child is built once here and both branches lower + // directly onto that build. A structured (array of objects, or a single object) input is + // exploded with mvexpand and each declared column is read with ITEM (typed ANY, since element + // types are erased upstream). A text input runs the split pipeline below on the same build. RelNode probe = node.getChild().get(0).accept(this, context); RelDataTypeField probeField = probe.getRowType().getField(node.getInField(), true, false); boolean structuredArray = @@ -4732,9 +4730,11 @@ public RelNode visitMultikv(Multikv node, CalcitePlanContext context) { if (noFields) { return relBuilder.peek(); } + List fieldTypes = node.getFieldTypes(); List projected = new ArrayList<>(); List names = new ArrayList<>(); - for (Field f : fields) { + for (int i = 0; i < fields.size(); i++) { + Field f = fields.get(i); String col = f.getField().toString(); RexNode item = PPLFuncImpTable.INSTANCE.resolve( @@ -4745,6 +4745,12 @@ public RelNode visitMultikv(Multikv node, CalcitePlanContext context) { col, context.rexBuilder.getTypeFactory().createSqlType(SqlTypeName.VARCHAR), true)); + ExprCoreType declared = fieldTypes == null ? null : fieldTypes.get(i); + if (declared != null) { + item = + context.rexBuilder.makeCast( + OpenSearchTypeFactory.convertExprTypeToRelDataType(declared), item, true, true); + } projected.add(item); names.add(col); } @@ -4752,10 +4758,32 @@ public RelNode visitMultikv(Multikv node, CalcitePlanContext context) { 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); - + // Reject input types multikv cannot read. The structured branch above handles array/object + // fields and text mode splits string values, so a scalar numeric/boolean/date field has no + // table text to parse; reject it at plan time instead of silently treating it as text. An + // untyped ANY field is allowed, since its runtime value may be text. + if (probeField != null) { + SqlTypeName inputTypeName = probeField.getType().getSqlTypeName(); + boolean textLike = inputTypeName == SqlTypeName.VARCHAR || inputTypeName == SqlTypeName.CHAR; + boolean untyped = inputTypeName == SqlTypeName.ANY || inputTypeName == SqlTypeName.NULL; + if (!textLike && !untyped) { + throw new SemanticCheckException( + "multikv input field '" + + node.getInField() + + "' has type " + + inputTypeName + + "; multikv reads table-formatted text or an array/object field. Cast it to a" + + " string first, for example: eval " + + node.getInField() + + " = cast(" + + node.getInField() + + " as string)."); + } + } + // Text input: lower the split pipeline directly onto the probe build above, so the child is + // visited exactly once. eval __multikv_record__ = MULTIKV_SPLIT(inField, ...), explode it with + // mvexpand (one row per table data row), then read the declared columns with MULTIKV_EXTRACT. + RelBuilder relBuilder = context.relBuilder; final String lineField = "__multikv_record__"; final int forceHeader = node.getForceHeader() == null ? -1 : node.getForceHeader(); final String filterJoined = @@ -4763,43 +4791,54 @@ public RelNode visitMultikv(Multikv node, CalcitePlanContext context) { ? "" : String.join(MultikvParser.FS, node.getFilterTerms()); - 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)))); + RexNode split = + PPLFuncImpTable.INSTANCE.resolve( + context.rexBuilder, + BuiltinFunctionName.MULTIKV_SPLIT, + relBuilder.field(node.getInField()), + relBuilder.literal(forceHeader), + relBuilder.literal(node.isNoHeader()), + context.rexBuilder.makeLiteral( + filterJoined, + context.rexBuilder.getTypeFactory().createSqlType(SqlTypeName.VARCHAR), + true)); + relBuilder.projectPlus(relBuilder.alias(split, lineField)); - plan = new MvExpand(AstDSL.field(lineField), null).attach(plan); + // mvexpand the record column: 1 -> N rows. Mirrors visitMvExpand (alias == field name). + buildExpandRelNode(relBuilder.field(lineField), lineField, lineField, null, context); 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); + return relBuilder.peek(); } - 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); - - UnresolvedExpression[] projections = fields.toArray(new UnresolvedExpression[0]); - plan = AstDSL.project(plan, projections); - - return plan.accept(this, context); + List fieldTypes = node.getFieldTypes(); + List projected = new ArrayList<>(); + List names = new ArrayList<>(); + for (int i = 0; i < fields.size(); i++) { + String col = fields.get(i).getField().toString(); + RexNode extract = + PPLFuncImpTable.INSTANCE.resolve( + context.rexBuilder, + BuiltinFunctionName.MULTIKV_EXTRACT, + relBuilder.field(lineField), + context.rexBuilder.makeLiteral( + col, + context.rexBuilder.getTypeFactory().createSqlType(SqlTypeName.VARCHAR), + true)); + ExprCoreType declared = fieldTypes == null ? null : fieldTypes.get(i); + if (declared != null) { + extract = + context.rexBuilder.makeCast( + OpenSearchTypeFactory.convertExprTypeToRelDataType(declared), extract, true, true); + } + projected.add(extract); + names.add(col); + } + relBuilder.project(projected, names); + context.setProjectVisited(true); + return relBuilder.peek(); } @Override diff --git a/docs/category.json b/docs/category.json index a8665b82eac..803948635b0 100644 --- a/docs/category.json +++ b/docs/category.json @@ -76,7 +76,8 @@ "user/ppl/cmd/table.md", "user/ppl/cmd/convert.md", "user/ppl/cmd/expand.md", - "user/ppl/cmd/flatten.md" + "user/ppl/cmd/flatten.md", + "user/ppl/cmd/multikv.md" ], "sql_cli": [ "user/dql/expressions.rst", diff --git a/docs/user/ppl/cmd/multikv.md b/docs/user/ppl/cmd/multikv.md index a42f1a0c9b3..c770e44c270 100644 --- a/docs/user/ppl/cmd/multikv.md +++ b/docs/user/ppl/cmd/multikv.md @@ -10,7 +10,7 @@ The `multikv` command extracts field values from an input field and emits one ro The `multikv` command has the following syntax: ```syntax -multikv [field=] [fields ...] [forceheader=] [noheader=] +multikv [field=] [fields [:]...] [forceheader=] [noheader=] ``` ## Parameters @@ -18,7 +18,7 @@ multikv [field=] [fields ...] [forceheader=] [noheader=] | 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`. | +| `fields [:]...` | Optional | Declares the output columns to extract by name. This is the only form that yields named columns. Each extracted column is `string` by default; append `:` (one of `string`, `boolean`, `int`/`integer`, `long`, `float`, `double`) to pin a type at plan time. | | `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`. | @@ -28,70 +28,115 @@ Every extracted column is typed as `string`, because the values come from splitt For structured input (an object or an array of objects), the extracted columns are typed `ANY` rather than the field's mapped scalar type, because object and nested fields collapse to `ANY`-valued containers before the command runs, so the mapped type is not recovered. Cast downstream for typed operations, for example `eval p = cast(pid as int)`. +To fix a column's type at plan time, append `:` in the `fields` clause, for example `fields pctIdle:long cpu:double`. The accepted types are `string`, `boolean`, `int`/`integer`, `long`, `float`, and `double`. This lowers to a safe cast, so a value that fails to parse arrives as null. It applies to both text and structured input and yields a typed column directly, so `fields pid:long` is equivalent to `fields pid` followed by `eval pid = cast(pid as long)`. A typed column name must begin with a letter or `*`. For UDT types (`date`, `time`, `timestamp`, `ip`, `json`), use `string` plus a downstream `cast`. + ## 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 +source=multikv_text | multikv field=raw fields pctIdle | fields pctIdle ``` -The query returns one row per table data row, with a single `pctIdle` (string) column. +```text +fetched rows / total rows = 2/2 ++---------+ +| pctIdle | +|---------| +| 90 | +| 92 | ++---------+ +``` ## Example 2: Extract multiple columns ```ppl -source=metrics +source=multikv_text | eval _raw = raw | multikv fields CPU pctIdle ``` -The query returns one row per table data row, with `CPU` (string) and `pctIdle` (string) columns. +```text +fetched rows / total rows = 2/2 ++-----+---------+ +| CPU | pctIdle | +|-----+---------| +| all | 90 | +| 0 | 92 | ++-----+---------+ +``` ## 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 +source=multikv_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. +```text +fetched rows / total rows = 2/2 ++----------+ +| endpoint | +|----------| +| foo | +| bar | ++----------+ +``` ## Example 4: Row explosion with noheader When there is no header row and only the number of data lines matters: ```ppl -source=lines +source=multikv_lines | eval _raw = raw | multikv noheader=true | stats count ``` -The query explodes the table text into one row per data line and counts them. +```text +fetched rows / total rows = 1/1 ++-------+ +| count | +|-------| +| 3 | ++-------+ +``` ## 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: +When `field=` points at an array of objects, `multikv` emits one row per element and reads each declared column from the element. The extracted columns are typed `ANY`, so cast downstream for typed operations: ```ppl -source=hosts +source=multikv_struct | 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=`. +```text +fetched rows / total rows = 2/2 ++-----+-----+ +| pid | cpu | +|-----+-----| +| 1 | 0.5 | +| 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 +```ppl ignore +source=multikv_text | 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. + +Extracted columns are untyped by default (`string` for text input, `ANY` for structured input). Typed operations therefore depend on the engine's per-operation coercion, an explicit downstream `cast`, or the inline `fields :` form (see [Column typing](#column-typing)). diff --git a/docs/user/ppl/index.md b/docs/user/ppl/index.md index 11b7651729e..26024812aea 100644 --- a/docs/user/ppl/index.md +++ b/docs/user/ppl/index.md @@ -84,7 +84,7 @@ source=accounts | [explain command](cmd/explain.md) | 3.1 | stable (since 3.1) | N/A | Explain the plan of query. | | [show datasources command](cmd/showdatasources.md) | 2.4 | stable (since 2.4) | N/A | Query datasources configured in the PPL engine. | | [makeresults command](cmd/makeresults.md) | 3.8 | experimental (since 3.8) | No | 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) | No | Extract fields from table-formatted text in a field, emitting one row per table data row. | +| [multikv command](cmd/multikv.md) | 3.9 | experimental (since 3.9) | No | 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) | Yes | Adds row and column values and appends a totals column and row. | | [addcoltotals command](cmd/addcoltotals.md) | 3.5 | stable (since 3.5) | Yes | Adds column values and appends a totals row. | | [transpose command](cmd/transpose.md) | 3.5 | stable (since 3.5) | Yes | Transpose rows to columns. | diff --git a/doctest/test_data/multikv_lines.json b/doctest/test_data/multikv_lines.json new file mode 100644 index 00000000000..171df9fce3f --- /dev/null +++ b/doctest/test_data/multikv_lines.json @@ -0,0 +1 @@ +{"raw": "a 1\nb 2\nc 3"} diff --git a/doctest/test_data/multikv_report.json b/doctest/test_data/multikv_report.json new file mode 100644 index 00000000000..4ed789d6440 --- /dev/null +++ b/doctest/test_data/multikv_report.json @@ -0,0 +1 @@ +{"raw": "== report ==\nendpoint count\nfoo 1\nbar 2"} diff --git a/doctest/test_data/multikv_struct.json b/doctest/test_data/multikv_struct.json new file mode 100644 index 00000000000..0782a0a7321 --- /dev/null +++ b/doctest/test_data/multikv_struct.json @@ -0,0 +1 @@ +{"procs":[{"pid":1,"cpu":0.5},{"pid":42,"cpu":9.1}]} diff --git a/doctest/test_data/multikv_text.json b/doctest/test_data/multikv_text.json new file mode 100644 index 00000000000..0217e70f137 --- /dev/null +++ b/doctest/test_data/multikv_text.json @@ -0,0 +1 @@ +{"raw": "CPU pctUser pctIdle\nall 5 90\n0 3 92"} diff --git a/doctest/test_docs.py b/doctest/test_docs.py index e179c85eb54..a8ce8eb7c54 100644 --- a/doctest/test_docs.py +++ b/doctest/test_docs.py @@ -60,6 +60,10 @@ 'time_test': 'time_test.json', 'mvcombine_data': 'mvcombine.json', 'timewrap_test': 'timewrap_test.json', + 'multikv_text': 'multikv_text.json', + 'multikv_report': 'multikv_report.json', + 'multikv_lines': 'multikv_lines.json', + 'multikv_struct': 'multikv_struct.json', } DEBUG_MODE = os.environ.get('DOCTEST_DEBUG', 'false').lower() == 'true' diff --git a/doctest/test_mapping/multikv_struct.json b/doctest/test_mapping/multikv_struct.json new file mode 100644 index 00000000000..dee394a597d --- /dev/null +++ b/doctest/test_mapping/multikv_struct.json @@ -0,0 +1 @@ +{"mappings":{"properties":{"procs":{"type":"nested","properties":{"pid":{"type":"long"},"cpu":{"type":"double"}}}}}} diff --git a/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalciteExplainIT.java b/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalciteExplainIT.java index d3589fcf138..a3d161a394b 100644 --- a/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalciteExplainIT.java +++ b/integ-test/src/test/java/org/opensearch/sql/calcite/remote/CalciteExplainIT.java @@ -64,6 +64,25 @@ public void init() throws Exception { loadIndex(Index.CASCADED_NESTED); loadIndex(Index.MVEXPAND_EDGE_CASES); loadIndex(Index.GRAPH_EMPLOYEES); + + // multikv fixtures: a text-table field and a structured array-of-objects field so the + // explain tests can plan both lowering paths. + if (!org.opensearch.sql.legacy.TestUtils.isIndexExist(client(), "test_multikv")) { + org.opensearch.sql.legacy.TestUtils.createIndexByRestClient(client(), "test_multikv", null); + Request doc = new Request("PUT", "/test_multikv/_doc/1?refresh=true"); + doc.setJsonEntity("{\"raw\": \"CPU pctUser pctIdle\\nall 5 90\\n0 3 92\"}"); + client().performRequest(doc); + } + if (!org.opensearch.sql.legacy.TestUtils.isIndexExist(client(), "test_multikv_struct")) { + String multikvStructMapping = + "{\"mappings\":{\"properties\":{\"procs\":{\"type\":\"nested\"," + + "\"properties\":{\"pid\":{\"type\":\"long\"},\"cpu\":{\"type\":\"double\"}}}}}}"; + org.opensearch.sql.legacy.TestUtils.createIndexByRestClient( + client(), "test_multikv_struct", multikvStructMapping); + 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); + } } // Only for Calcite: the rest row source explains as a CalciteScannableCatalogScan. @@ -79,6 +98,28 @@ public void explainRestCommand() throws IOException { @Ignore("test only in v2") public void testExplainModeUnsupportedInV2() throws IOException {} + // multikv text mode lowers to a MULTIKV_EXTRACT projection over the source field. + @Test + public void explainMultikvTextMode() throws IOException { + String result = + explainQueryToString( + "source=test_multikv | multikv field=raw fields pctIdle | fields pctIdle"); + Assert.assertTrue( + "Expected MULTIKV_EXTRACT in the explain output, got: " + result, + result.contains("MULTIKV_EXTRACT")); + } + + // fields col:type wraps the extracted column in a plan-time safe cast to the declared type. + @Test + public void explainMultikvTypedTextColumn() throws IOException { + String result = + explainQueryToString( + "source=test_multikv | multikv field=raw fields pctIdle:long | fields pctIdle"); + Assert.assertTrue( + "Expected a SAFE_CAST over MULTIKV_EXTRACT in the explain output, got: " + result, + result.contains("SAFE_CAST") && result.contains("MULTIKV_EXTRACT")); + } + // Only for Calcite @Test public void supportSearchSargPushDown_singleRange() throws IOException { 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 index 7e06d15e1f2..e880ccbc6d3 100644 --- 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 @@ -139,6 +139,59 @@ public void testMultikvSingleObject() throws IOException { verifyDataRows(result, rows("root", "admin")); } + @Test + public void testMultikvTypedTextColumn() throws IOException { + // fields col:type types pctIdle as long at plan time, so it filters and returns as a number. + JSONObject result = + executeQuery( + "source=test_multikv | eval _raw = raw | multikv fields pctIdle:long" + + " | where pctIdle > 90 | fields pctIdle"); + verifySchema(result, schema("pctIdle", "bigint")); + verifyDataRows(result, rows(92)); + } + + @Test + public void testMultikvTypedStructuredColumns() throws IOException { + // Structured col:type casts each exploded column to the declared type. + JSONObject result = + executeQuery( + "source=test_multikv_struct | multikv field=procs fields pid:long cpu:double" + + " | where pid > 1 | fields pid cpu"); + verifySchema(result, schema("pid", "bigint"), schema("cpu", "double")); + verifyDataRows(result, rows(42, 9.1)); + } + + @Test + public void testMultikvTextCastThenFilterKeepsSiblingColumn() throws IOException { + // Downstream consume: cast one extracted column, filter on it, project a sibling column. + JSONObject result = + executeQuery( + "source=test_multikv | eval _raw = raw | multikv fields CPU pctIdle" + + " | where cast(pctIdle as int) > 90 | fields CPU"); + verifyDataRows(result, rows("0")); + } + + @Test + public void testMultikvTypedTextFeedsAggregate() throws IOException { + // A typed text column feeds a numeric aggregate with no explicit downstream cast. + JSONObject result = + executeQuery( + "source=test_multikv | eval _raw = raw | multikv fields pctIdle:long" + + " | stats sum(pctIdle) as total"); + verifyDataRows(result, rows(182)); + } + + @Test + public void testMultikvTypedStructuredSortAndProject() throws IOException { + // Typed structured columns sort numerically downstream with no explicit cast. + JSONObject result = + executeQuery( + "source=test_multikv_struct | multikv field=procs fields pid:long cpu:double" + + " | sort pid | fields pid cpu"); + verifySchema(result, schema("pid", "bigint"), schema("cpu", "double")); + verifyDataRows(result, rows(1, 0.5), rows(42, 9.1)); + } + @Test public void testBareMultikvRejectedWithGuidance() throws IOException { // Bare auto-header multikv referencing a column downstream must be rejected at plan time diff --git a/ppl/src/main/antlr/OpenSearchPPLParser.g4 b/ppl/src/main/antlr/OpenSearchPPLParser.g4 index 5d01668e009..1ecdb2f9f99 100644 --- a/ppl/src/main/antlr/OpenSearchPPLParser.g4 +++ b/ppl/src/main/antlr/OpenSearchPPLParser.g4 @@ -678,13 +678,25 @@ multikvCommand ; multikvParameter - : FIELDS fields = fieldList + : FIELDS fields = multikvFieldList | FIELD EQUAL inField = qualifiedName | FORCEHEADER EQUAL forceHeader = integerLiteral | NOHEADER EQUAL noHeader = booleanLiteral | RMORIG EQUAL rmOrig = booleanLiteral ; +// multikv-local typed field list (keeps the shared fieldList untouched). The case-insensitive lexer +// folds "col:" into one CLUSTER token, so the typed form is CLUSTER + type and the name is that +// token minus its trailing colon (a typed name must start with a letter or '*'). +multikvFieldList + : multikvField ((COMMA)? multikvField)* + ; + +multikvField + : CLUSTER convertedDataType + | fieldExpression + ; + 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 1c81ee42f7a..c282f157ef3 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 @@ -138,6 +138,7 @@ import org.opensearch.sql.common.setting.Settings; import org.opensearch.sql.common.setting.Settings.Key; import org.opensearch.sql.common.utils.StringUtils; +import org.opensearch.sql.data.type.ExprCoreType; import org.opensearch.sql.exception.SemanticCheckException; import org.opensearch.sql.ppl.antlr.parser.OpenSearchPPLParser; import org.opensearch.sql.ppl.antlr.parser.OpenSearchPPLParser.AdCommandContext; @@ -150,6 +151,7 @@ import org.opensearch.sql.ppl.antlr.parser.OpenSearchPPLParserBaseVisitor; import org.opensearch.sql.ppl.utils.ArgumentFactory; import org.opensearch.sql.ppl.utils.MakeResultsDataParser; +import org.opensearch.sql.ppl.utils.PplInlineTypeResolver; import org.opensearch.sql.ppl.utils.UnresolvedPlanHelper; import org.opensearch.sql.utils.SystemIndexUtils; @@ -1111,17 +1113,35 @@ public UnresolvedPlan visitMvexpandCommand(OpenSearchPPLParser.MvexpandCommandCo @Override public UnresolvedPlan visitMultikvCommand(OpenSearchPPLParser.MultikvCommandContext ctx) { List fields = null; + List fieldTypes = null; String inField = Multikv.DEFAULT_INPUT_FIELD; Integer forceHeader = null; boolean noHeader = false; - boolean rmOrig = true; // Splunk default + boolean rmOrig = true; // drop the original event fields unless keepOrig is set for (OpenSearchPPLParser.MultikvParameterContext p : ctx.multikvParameter()) { if (p.fields != null) { - fields = - p.fields.fieldExpression().stream() - .map(f -> (Field) expressionBuilder.visit(f)) - .collect(Collectors.toList()); + fields = new ArrayList<>(); + fieldTypes = new ArrayList<>(); + boolean anyTyped = false; + for (OpenSearchPPLParser.MultikvFieldContext mf : p.fields.multikvField()) { + if (mf.CLUSTER() != null) { + // The lexer folds "col:" into one CLUSTER token; drop the trailing colon for the name. + String raw = mf.CLUSTER().getText(); + String col = StringUtils.unquoteIdentifier(raw.substring(0, raw.length() - 1)); + fields.add(new Field(new QualifiedName(col))); + fieldTypes.add( + PplInlineTypeResolver.resolve( + mf.convertedDataType().typeName.getText(), "multikv")); + anyTyped = true; + } else { + fields.add((Field) expressionBuilder.visit(mf.fieldExpression())); + fieldTypes.add(null); + } + } + if (!anyTyped) { + fieldTypes = null; + } } else if (p.inField != null) { inField = p.inField.getText(); } else if (p.forceHeader != null) { @@ -1137,7 +1157,7 @@ public UnresolvedPlan visitMultikvCommand(OpenSearchPPLParser.MultikvCommandCont } } - return new Multikv(inField, fields, null, forceHeader, noHeader, rmOrig); + return new Multikv(inField, fields, fieldTypes, null, forceHeader, noHeader, rmOrig); } @Override diff --git a/ppl/src/main/java/org/opensearch/sql/ppl/utils/MakeResultsDataParser.java b/ppl/src/main/java/org/opensearch/sql/ppl/utils/MakeResultsDataParser.java index 225f3d0280f..cd69a737f0d 100644 --- a/ppl/src/main/java/org/opensearch/sql/ppl/utils/MakeResultsDataParser.java +++ b/ppl/src/main/java/org/opensearch/sql/ppl/utils/MakeResultsDataParser.java @@ -318,30 +318,7 @@ private static Object jsonValue(JsonNode v) { } private static ExprCoreType resolveType(String name) { - switch (name.toLowerCase(Locale.ROOT)) { - case "string": - return ExprCoreType.STRING; - case "boolean": - return ExprCoreType.BOOLEAN; - case "int": - case "integer": - return ExprCoreType.INTEGER; - case "long": - return ExprCoreType.LONG; - case "float": - return ExprCoreType.FLOAT; - case "double": - return ExprCoreType.DOUBLE; - case "date": - case "time": - case "timestamp": - case "ip": - case "json": - throw new SyntaxCheckException( - "makeresults inline type '" + name + "' is not yet supported; use string and cast"); - default: - return null; - } + return PplInlineTypeResolver.resolve(name, "makeresults"); } private static Object coerce(Object value, ExprCoreType type) { diff --git a/ppl/src/main/java/org/opensearch/sql/ppl/utils/PplInlineTypeResolver.java b/ppl/src/main/java/org/opensearch/sql/ppl/utils/PplInlineTypeResolver.java new file mode 100644 index 00000000000..38385e147d5 --- /dev/null +++ b/ppl/src/main/java/org/opensearch/sql/ppl/utils/PplInlineTypeResolver.java @@ -0,0 +1,53 @@ +/* + * Copyright OpenSearch Contributors + * SPDX-License-Identifier: Apache-2.0 + */ + +package org.opensearch.sql.ppl.utils; + +import java.util.Locale; +import org.opensearch.sql.common.antlr.SyntaxCheckException; +import org.opensearch.sql.data.type.ExprCoreType; + +/** + * Resolves an inline PPL scalar type name (the {@code makeresults data="name:type"} header and the + * {@code multikv fields col:type} clause) to its {@link ExprCoreType}. An unknown name returns + * {@code null} so the caller can fall back to string; a UDT name (date/time/timestamp/ip/json) is + * rejected, since the inline paths only lower to a cast over string values. + */ +public final class PplInlineTypeResolver { + + private PplInlineTypeResolver() {} + + /** + * @param commandName only used to phrase the rejection message + * @return the resolved type, or {@code null} when {@code name} is not a known scalar type + * @throws SyntaxCheckException when {@code name} is an unsupported UDT type + */ + public static ExprCoreType resolve(String name, String commandName) { + switch (name.toLowerCase(Locale.ROOT)) { + case "string": + return ExprCoreType.STRING; + case "boolean": + return ExprCoreType.BOOLEAN; + case "int": + case "integer": + return ExprCoreType.INTEGER; + case "long": + return ExprCoreType.LONG; + case "float": + return ExprCoreType.FLOAT; + case "double": + return ExprCoreType.DOUBLE; + case "date": + case "time": + case "timestamp": + case "ip": + case "json": + throw new SyntaxCheckException( + commandName + " inline type '" + name + "' is not yet supported; use string and cast"); + default: + return null; + } + } +} 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 index abe55d63822..5c4d3797a17 100644 --- a/ppl/src/test/java/org/opensearch/sql/ppl/calcite/CalcitePPLMultikvTest.java +++ b/ppl/src/test/java/org/opensearch/sql/ppl/calcite/CalcitePPLMultikvTest.java @@ -91,4 +91,29 @@ public void testBareMultikvRejectedWithGuidance() { // rejected at the field-resolution phase with actionable guidance. expectError("source=EMP | eval _raw = ENAME | multikv", "fields clause"); } + + @Test + public void testMultikvTypedTextColumnAddsCast() { + // A declared type lowers to a SAFE_CAST over the extraction; the untyped form has none. + String typed = + plan("source=EMP | eval _raw = ENAME | multikv fields pctIdle:long | fields pctIdle"); + assertTrue(typed, typed.contains("MULTIKV_EXTRACT")); + assertTrue(typed, typed.contains("SAFE_CAST")); + String untyped = + plan("source=EMP | eval _raw = ENAME | multikv fields pctIdle | fields pctIdle"); + assertFalse(untyped, untyped.contains("SAFE_CAST")); + } + + @Test + public void testMultikvRejectsUnsupportedInlineType() { + expectError( + "source=EMP | eval _raw = ENAME | multikv fields pctIdle:ip", "is not yet supported"); + } + + @Test + public void testMultikvRejectsNonTextInputField() { + // A scalar numeric input field (EMPNO) has no table text to parse, so multikv rejects it at + // plan time rather than silently treating it as text. + expectError("source=EMP | multikv field=EMPNO fields x", "reads table-formatted text"); + } } 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 8a02fb57e62..0e3a90b5b05 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 @@ -87,6 +87,7 @@ import org.opensearch.sql.ast.tree.Xyseries; import org.opensearch.sql.common.antlr.SyntaxCheckException; import org.opensearch.sql.common.setting.Settings.Key; +import org.opensearch.sql.data.type.ExprCoreType; import org.opensearch.sql.exception.SemanticCheckException; import org.opensearch.sql.ppl.AstPlanningTestBase; import org.opensearch.sql.utils.SystemIndexUtils; @@ -1225,6 +1226,33 @@ public void testMultikvCommand() { .attach(relation("t"))); } + @Test + public void testMultikvCommandWithTypedFields() { + assertEqual( + "source=t | multikv fields pctIdle:long", + new Multikv( + "_raw", + Arrays.asList(field("pctIdle")), + Arrays.asList(ExprCoreType.LONG), + null, + null, + false, + true) + .attach(relation("t"))); + // an undeclared column keeps a null slot alongside the typed ones + assertEqual( + "source=t | multikv field=message fields CPU pctIdle:long ratio:double", + new Multikv( + "message", + Arrays.asList(field("CPU"), field("pctIdle"), field("ratio")), + Arrays.asList(null, ExprCoreType.LONG, ExprCoreType.DOUBLE), + null, + null, + false, + true) + .attach(relation("t"))); + } + @Test public void testDescribeMatchAllCrossClusterSearchCommand() { assertEqual("describe *:t", describe(mappingTable("*:t")));