Skip to content

fix(analytics): Fix multi-shard GROUP BY on multi_value keys - #23091

Merged
linuxpi merged 4 commits into
opensearch-project:mainfrom
linuxpi:fix/mv-multishard-groupby-decode
Oct 1, 2026
Merged

linuxpi merged 4 commits into
opensearch-project:mainfrom
linuxpi:fix/mv-multishard-groupby-decode

Conversation

@linuxpi

@linuxpi linuxpi commented Sep 18, 2026

Copy link
Copy Markdown
Contributor

Description

Implicit GROUP BY on a multi_value keyword field (e.g. source = idx | stats count() by tags) failed on any index with more than one shard. Single-shard queries were unaffected because they never build a coordinator reduce fragment. Two layers on the reduce stage were wrong:

  1. Plan assembly (400). attachFragmentOnTop round-trips the shard (PARTIAL) fragment through the stock substrait-java ProtoPlanConverter, which maps the opensearch://analytics/multi_value_expand/v1 ExtensionSingle to an EmptyDetail with an empty record type. The grouping key appended by the expansion was then out of range: Field reference offset (N) must be less than number of fields in struct (0). MultiValueExpandDetail gains a fromProto inverse and derives its record type from the input (append or replace the LIST column with its nullable element), and decodePlan uses a ProtoPlanConverter whose detailFromExtensionSingleRel recognizes the extension type URL.

  2. Execution (500). The FINAL aggregate's StageInputTableScan still carried the pre-expansion LIST type while the shards stream the already-expanded scalar key, so DataFusion rejected the reduce ReadRel: Field 'tags' in Substrait schema has a different type (List(Utf8)) than the corresponding field in the table schema (Utf8View). OpenSearchAggregate.stripAnnotations now tags FINAL-mode aggregates with a RelHint, and MultiValueRelRewriter retypes LIST group keys on a FINAL aggregate over a stage input to their element type instead of re-expanding them; PARTIAL aggregates keep the existing expansion path.

Tests: DataFusionFragmentConvertorTests adds a decode round-trip test (reproduces the offset error on old code) and a FINAL-retype test (no expand extension emitted, key declared as scalar in the ReadRel schema). The three MultiValueAggregationIT cases pinned on #23057 (count() by tags, sum(latency) by tags, region, SQL GROUP BY tags, region on 2 shards) are un-skipped and pass against the live 2-node cluster.

Supersedes #23067 with a minimal backend-only change (no planner refactor).

Related Issues

Resolves #23057

Check List

  • Functionality includes testing.
  • API changes companion pull request created, if applicable.
  • Public documentation issue/PR created, if applicable.

By submitting this pull request, I confirm that my contribution is made under the terms of the Apache 2.0 license.
For more information on following Developer Certificate of Origin and signing off your commits, please check here.

@linuxpi
linuxpi requested a review from a team as a code owner September 18, 2026 21:07
@github-actions github-actions Bot added bug Something isn't working Search Search query, autocomplete ...etc labels Sep 18, 2026
@github-actions

Copy link
Copy Markdown
Contributor

PR Code Analyzer ❗

AI-powered 'Code-Diff-Analyzer' found issues on commit d7e9e3d.

⛔ Hard block: Issues at Medium severity or above will block this PR from merging.

PathLineSeverityDescription
sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/DataFusionFragmentConvertor.java1204lowDebug log statement `LOGGER.info("---------i was here!!")` left in production code path inside `decodePlan`. The dashed prefix and casual phrasing are hallmarks of temporary debugging code accidentally committed. This leaks internal execution flow information to logs but has no malicious indicators.

The table above displays the top 10 most important findings.

Total: 1 | Critical: 0 | High: 0 | Medium: 0 | Low: 1


Pull Requests Author(s): Please update your Pull Request according to the report above.

Repository Maintainer(s): You can bypass diff analyzer by adding label skip-diff-analyzer after reviewing the changes carefully, then re-run failed actions. To re-enable the analyzer, remove the label, then re-run all actions.


⚠️ Note: The Code-Diff-Analyzer helps protect against potentially harmful code patterns. Please ensure you have thoroughly reviewed the changes beforehand.

Thanks.

@github-actions

github-actions Bot commented Sep 18, 2026 •

Copy link
Copy Markdown
Contributor

PR Reviewer Guide 🔍

(Review updated until commit bc83ce6)

Here are some key observations to aid the review process:

🧪 PR contains tests
🔒 No security concerns identified
✅ No TODO sections
🔀 No multiple PR themes
⚡ No major issues detected

@github-actions

github-actions Bot commented Sep 18, 2026 •

Copy link
Copy Markdown
Contributor

PR Code Suggestions ✨

Latest suggestions up to bc83ce6

Explore these optional code suggestions:

CategorySuggestion                                                                                                                                    Impact
Possible issue
Enforce explicit byte order on decode

The payload is written with ByteBuffer.allocate(16) (default big-endian) but the
reader uses whatever byte order asReadOnlyByteBuffer() returns from a ByteString. To
guarantee round-trip compatibility across JVMs/platforms, explicitly set
ByteOrder.BIG_ENDIAN on both the writer and reader buffers.

sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/DataFusionFragmentConvertor.java [939-943]

-ByteBuffer payload = value.asReadOnlyByteBuffer();
+ByteBuffer payload = value.asReadOnlyByteBuffer().order(java.nio.ByteOrder.BIG_ENDIAN);
 int fieldIndex = payload.getInt();
 int limit = payload.getInt();
 int append = payload.getInt();
 int distinct = payload.getInt();
Suggestion importance[1-10]: 5

__

Why: ByteBuffer.allocate defaults to big-endian, and ByteString.asReadOnlyByteBuffer() also returns a big-endian buffer, so the round-trip works today. Still, explicitly setting the byte order is a small defensive improvement for cross-platform clarity.

Low
General
Preserve full qualified table name

Using getQualifiedName().getFirst() drops any schema/catalog qualifiers if the table
name is multi-part, which may produce a stage input pointing at a different (or
missing) table. Preserve the full qualified name (or reconstruct with the same
qualifier the original scan used) to avoid breaking the FINAL fragment lookup.

sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/MultiValueRelRewriter.java [128-133]

 RelNode retyped = new DataFusionFragmentConvertor.StageInputTableScan(
     stageInput.getCluster(),
     stageInput.getTraitSet(),
-    stageInput.getTable().getQualifiedName().getFirst(),
+    String.join(".", stageInput.getTable().getQualifiedName()),
     builder.build()
 );
Suggestion importance[1-10]: 3

__

Why: The StageInputTableScan constructor takes a single string identifier (e.g. "input-1"), so using getFirst() matches the existing convention. The suggested change to join qualifiers could actually break lookups rather than fix them.

Low

Previous suggestions

Suggestions up to commit ae1d4a4
CategorySuggestion                                                                                                                                    Impact
Possible issue
Enforce explicit byte order on decode

The writer side uses ByteBuffer.allocate(16) with default big-endian order, but
Any.getValue().asReadOnlyByteBuffer() may return a buffer whose byte order is not
guaranteed to match. Explicitly set ByteOrder.BIG_ENDIAN on both sides to guarantee
round-trip correctness across JVMs and platforms.

sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/DataFusionFragmentConvertor.java [939-943]

-ByteBuffer payload = value.asReadOnlyByteBuffer();
+ByteBuffer payload = value.asReadOnlyByteBuffer().order(java.nio.ByteOrder.BIG_ENDIAN);
 int fieldIndex = payload.getInt();
 int limit = payload.getInt();
 int append = payload.getInt();
 int distinct = payload.getInt();
Suggestion importance[1-10]: 3

__

Why: While ByteBuffer.allocate defaults to big-endian, so does asReadOnlyByteBuffer() per the ByteString contract; explicitly setting the order is a minor defensive improvement rather than a correctness fix.

Low
Suggestions up to commit be771bd
CategorySuggestion                                                                                                                                    Impact
Possible issue
Preserve full qualified name when retyping

getQualifiedName().getFirst() drops all but the first name segment of the stage
input table. If the qualified name has multiple segments (schema/table), the retyped
scan will point to the wrong table. Preserve the full qualified name (or reuse the
original table reference) when constructing the retyped scan.

sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/MultiValueRelRewriter.java [128-133]

 RelNode retyped = new DataFusionFragmentConvertor.StageInputTableScan(
     stageInput.getCluster(),
     stageInput.getTraitSet(),
-    stageInput.getTable().getQualifiedName().getFirst(),
+    String.join(".", stageInput.getTable().getQualifiedName()),
     builder.build()
 );
Suggestion importance[1-10]: 5

__

Why: Using getFirst() may drop qualified name segments if the table name has multiple parts, but the suggested fix (joining with ".") may not be semantically correct either; the concern is valid but depends on the StageInputTableScan constructor contract.

Low
General
Pin byte order for wire payload decoding

ByteBuffer.asReadOnlyByteBuffer() returns a buffer whose byte order defaults to
big-endian, which happens to match the writer here, but the write side uses
ByteBuffer.allocate(PAYLOAD_BYTES) without explicitly setting order either. To make
the wire format explicit and robust against JVM/library defaults, set
ByteOrder.BIG_ENDIAN on both the read and write buffers so decode never silently
misinterprets the payload.

sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/DataFusionFragmentConvertor.java [939-943]

-ByteBuffer payload = value.asReadOnlyByteBuffer();
+ByteBuffer payload = value.asReadOnlyByteBuffer().order(java.nio.ByteOrder.BIG_ENDIAN);
 int fieldIndex = payload.getInt();
 int limit = payload.getInt();
 int append = payload.getInt();
 int distinct = payload.getInt();
Suggestion importance[1-10]: 4

__

Why: ByteBuffer defaults to big-endian, so this is not a correctness bug today, but explicitly setting the byte order improves robustness and clarity of the wire format contract.

Low
Suggestions up to commit f7b202d
CategorySuggestion                                                                                                                                    Impact
General
Guard retyping against aggregate-call argument overlap

The retyping unconditionally forces the element type to nullable, but it does not
check whether the aggregate has aggregate calls that reference this field index as
an argument. If a non-grouping aggregate call references a LIST column, only
grouping keys are retyped here — that is intended — but any aggregate call whose
argument index points at the retyped field will now see a scalar element type
instead of the expected LIST, potentially causing a signature mismatch. Verify that
no AggregateCall argument list overlaps with aggregate.getGroupSet(), or restrict
retyping to indices that are exclusively grouping keys.

sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/MultiValueRelRewriter.java [117-123]

 RelDataType elementType = field.getType().getComponentType();
 if (elementType != null && aggregate.getGroupSet().get(fieldIndex)) {
+    // Grouping keys are never also aggregate-call arguments in a well-formed Aggregate,
+    // but assert this to catch shape regressions early.
+    assert aggregate.getAggCallList().stream().noneMatch(c -> c.getArgList().contains(fieldIndex));
     builder.add(field.getName(), typeFactory.createTypeWithNullability(elementType, true));
     changed = true;
 } else {
     builder.add(field.getName(), field.getType());
 }
Suggestion importance[1-10]: 3

__

Why: Adding a defensive assertion is a minor code hygiene improvement, but in a well-formed Calcite Aggregate, grouping keys and aggregate call arguments don't conflict in the way described, so the practical impact is low.

Low
Explicitly set big-endian byte order on decode

Any.getValue().asReadOnlyByteBuffer() returns a buffer whose byte order defaults to
big-endian, matching ByteBuffer.allocate(16) on the write side, but relying on this
implicit default is fragile — especially since the Rust side (per the comment) is
the peer producer. Explicitly set ByteOrder.BIG_ENDIAN on both the read and write
buffers so the wire contract is documented in code and immune to any future JVM
changes or refactors that swap the buffer implementation.

sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/DataFusionFragmentConvertor.java [939-943]

-ByteBuffer payload = value.asReadOnlyByteBuffer();
+ByteBuffer payload = value.asReadOnlyByteBuffer().order(java.nio.ByteOrder.BIG_ENDIAN);
 int fieldIndex = payload.getInt();
 int limit = payload.getInt();
 int append = payload.getInt();
 int distinct = payload.getInt();
Suggestion importance[1-10]: 3

__

Why: Explicitly setting ByteOrder.BIG_ENDIAN documents the wire contract and hardens against future refactors, but ByteBuffer is guaranteed big-endian by default in Java, so the correctness impact is minimal.

Low
Set explicit byte order for wire payload

Any.getTypeUrl() returns an empty string, not null, when unset, but comparing with
TYPE_URL.equals(...) handles both safely; however, the outer null check on any is
defensive but the caller in detailFromExtensionSingleRel never passes null. More
importantly, fromProto reads the payload with the default big-endian ByteBuffer
order but does not explicitly set it — while big-endian is the JVM default, the
write side also relies on this implicit default. Set the byte order explicitly on
both sides to guard against future refactors introducing a mismatch that would
silently corrupt the wire format across the Rust boundary.

sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/DataFusionFragmentConvertor.java [924-926]

 static boolean matches(Any any) {
     return any != null && TYPE_URL.equals(any.getTypeUrl());
 }
+// Note: ensure ByteBuffer order is explicitly BIG_ENDIAN in both fromProto and toProto.
Suggestion importance[1-10]: 2

__

Why: The suggestion is largely a note in a comment; the improved_code is essentially identical to existing_code with only a comment added. While explicit byte order can be a good practice, ByteBuffer defaults to big-endian per Java spec, making this a marginal improvement.

Low
Suggestions up to commit d7e9e3d
CategorySuggestion                                                                                                                                    Impact
General
Remove leftover debug log statement

Remove the leftover debug log statement LOGGER.info("---------i was here!!");. It
looks like a scratch trace, will pollute logs on every plan decode (a hot path on
the coordinator), and provides no diagnostic value.

sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/DataFusionFragmentConvertor.java [1206]

 private Plan decodePlan(byte[] bytes) {
-    LOGGER.info("---------i was here!!");
     try {
         io.substrait.proto.Plan proto = io.substrait.proto.Plan.parseFrom(bytes);
         return new OpenSearchProtoPlanConverter(extensions).from(proto);
Suggestion importance[1-10]: 9

__

Why: The LOGGER.info("---------i was here!!") is clearly leftover debug output that should not be shipped; removing it prevents log pollution on a hot coordinator path.

High
Pin byte order for wire payload decoding

The write side uses default (big-endian) byte order but reads should be explicit to
guarantee symmetry. ByteString.asReadOnlyByteBuffer() returns a buffer whose order
defaults to big-endian, matching writes, but explicitly setting ByteOrder.BIG_ENDIAN
documents the wire contract and protects against future writer changes.

sandbox/plugins/analytics-backend-datafusion/src/main/java/org/opensearch/be/datafusion/DataFusionFragmentConvertor.java [939-943]

-ByteBuffer payload = value.asReadOnlyByteBuffer();
+ByteBuffer payload = value.asReadOnlyByteBuffer().order(java.nio.ByteOrder.BIG_ENDIAN);
 int fieldIndex = payload.getInt();
 int limit = payload.getInt();
 int append = payload.getInt();
 int distinct = payload.getInt();
Suggestion importance[1-10]: 3

__

Why: ByteBuffer defaults to big-endian so this is a minor documentation/robustness improvement rather than a functional fix.

Low

Implicit GROUP BY on a multi_value keyword field failed on any index
with more than one shard, in two layers on the coordinator reduce
stage.

First, attachFragmentOnTop round-trips the shard (PARTIAL) fragment
through the stock substrait-java ProtoPlanConverter, which maps the
opensearch://analytics/multi_value_expand/v1 ExtensionSingle to an
EmptyDetail with an empty record type. The grouping key that the
expansion appended was then out of range, failing plan assembly with
"Field reference offset (N) must be less than number of fields in
struct (0)". MultiValueExpandDetail gains a fromProto inverse and
derives its record type from the input (append or replace the LIST
column with its nullable element), and decodePlan uses a
ProtoPlanConverter whose detailFromExtensionSingleRel recognizes the
extension type URL.

Second, the FINAL aggregate's StageInputTableScan still carried the
pre-expansion LIST type while the shards stream the already-expanded
scalar key, so DataFusion rejected the reduce ReadRel with "Field
'tags' in Substrait schema has a different type (List(Utf8)) than the
corresponding field in the table schema (Utf8View)". OpenSearchAggregate
now tags FINAL-mode aggregates with a RelHint when stripping
annotations, and MultiValueRelRewriter retypes LIST group keys on a
FINAL aggregate over a stage input to their element type instead of
re-expanding them.

Unit tests cover the decode round trip and the FINAL retype path; the
three MultiValueAggregationIT cases pinned on the issue are un-skipped.

Resolves opensearch-project#23057

Signed-off-by: Varun Bansal <bansvaru@amazon.com>
@linuxpi
linuxpi force-pushed the fix/mv-multishard-groupby-decode branch from d7e9e3d to f7b202d Compare September 18, 2026 21:10
@github-actions

Copy link
Copy Markdown
Contributor

Persistent review updated to latest commit f7b202d

@github-actions

Copy link
Copy Markdown
Contributor

❌ Gradle check result for f7b202d: FAILURE

Please examine the workflow log, locate, and copy-paste the failure(s) below, then iterate to green. Is the failure a flaky test unrelated to your change?

@github-actions

Copy link
Copy Markdown
Contributor

Persistent review updated to latest commit be771bd

@github-actions

Copy link
Copy Markdown
Contributor

✅ Gradle check result for be771bd: SUCCESS

@codecov

codecov Bot commented Sep 22, 2026 •

Copy link
Copy Markdown

Codecov Report

✅ All modified and coverable lines are covered by tests.
✅ Project coverage is 71.78%. Comparing base (afd64ec) to head (bc83ce6).
⚠️ Report is 11 commits behind head on main.

Additional details and impacted files
@@             Coverage Diff              @@
##               main   #23091      +/-   ##
============================================
- Coverage     71.78%   71.78%   -0.01%     
- Complexity    77773    77792      +19     
============================================
  Files          6179     6179              
  Lines        360740   360749       +9     
  Branches      52506    52507       +1     
============================================
- Hits         258971   258960      -11     
- Misses        81172    81179       +7     
- Partials      20597    20610      +13     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

@github-actions

Copy link
Copy Markdown
Contributor

Persistent review updated to latest commit ae1d4a4

@github-actions

Copy link
Copy Markdown
Contributor

❌ Gradle check result for ae1d4a4: FAILURE

Please examine the workflow log, locate, and copy-paste the failure(s) below, then iterate to green. Is the failure a flaky test unrelated to your change?

@bharath-techie bharath-techie left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

@github-actions

Copy link
Copy Markdown
Contributor

Persistent review updated to latest commit bc83ce6

@github-actions

Copy link
Copy Markdown
Contributor

❌ Gradle check result for bc83ce6: FAILURE

Please examine the workflow log, locate, and copy-paste the failure(s) below, then iterate to green. Is the failure a flaky test unrelated to your change?

@github-actions

Copy link
Copy Markdown
Contributor

❕ Gradle check result for bc83ce6: UNSTABLE

Please review all flaky tests that succeeded after retry and create an issue if one does not already exist to track the flaky failure.

@linuxpi
linuxpi merged commit 4247b3d into opensearch-project:main Oct 1, 2026
19 of 25 checks passed
@linuxpi
linuxpi deleted the fix/mv-multishard-groupby-decode branch October 1, 2026 08:58
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working Search Search query, autocomplete ...etc

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[BUG][Sandbox] Multi-shard GROUP BY on a multi_value keyword fails with 'Field reference offset must be less than number of fields in struct (0)'

3 participants