Convert dedup pushdown to composite + top_hits - #4844

Merged
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797
Nov 26, 2025
Merged

Convert dedup pushdown to composite + top_hits#4844
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797

Conversation

@LantaoJin

@LantaoJinLantaoJin commented Nov 21, 2025

Copy link
Copy Markdown
Member

Description

Convert dedup pushdown to composite + top_hits
Main changes:

  1. Add a rule DedupPushdownRule to convert the dedup plan pattern to Aggregate with TopHits metrics
  2. Upgrade TopHitsParser to support parsing data from _source (fetchField cannot work for OS Object type)
  3. Corresponding API change on MetricParser: Map<String, Object> parse(Aggregation aggregation); -> List<Map<String, Object>> parse(Aggregation aggregation);

Follow-ups: Support dedup on script (expression)

Related Issues

Resolves#4797

Check List

  • New functionality includes testing.
  • New functionality has been documented.
  • New functionality has javadoc added.
  • New functionality has a user manual doc added.
  • New PPL command checklist all confirmed.
  • API changes companion pull request created.
  • Commits are signed per the DCO using --signoff or -s.
  • Public documentation issue/PR created.

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.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@LantaoJinLantaoJin added the enhancement New feature or request label Nov 21, 2025
@LantaoJinLantaoJin mentioned this pull request Nov 21, 2025
8 tasks
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
return project == null
? List.of()
: aggCall.getArgList().stream().map(project.getProjects()::get).toList();
: PlanUtils.getObjectFromLiteralAgg(aggCall) != null

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why calling this method here?

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Is it for identifying LITERAL_AGG?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. LITERAL_AGG passes the numberOfDedup here

LogicalFilter(condition=[IS NOT NULL($4)])
CalciteLogicalIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]])
physical: |
CalciteEnumerableIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]], PushDownContext=[[PROJECT->[account_number, firstname, address, balance, gender, city, employer, state, age, email, lastname], FILTER->IS NOT NULL($4), AGGREGATION->rel#:LogicalAggregate.NONE.[](input=LogicalProject#,group={0},agg#0=LITERAL_AGG(1)), LIMIT->10000], OpenSearchRequestBuilder(sourceBuilder={"from":0,"size":0,"timeout":"1m","query":{"exists":{"field":"gender","boost":1.0}},"_source":{"includes":["account_number","firstname","address","balance","gender","city","employer","state","age","email","lastname"],"excludes":[]},"aggregations":{"composite_buckets":{"composite":{"size":1000,"sources":[{"gender":{"terms":{"field":"gender.keyword","missing_bucket":false,"order":"asc"}}}]},"aggregations":{"$f1":{"top_hits":{"from":0,"size":1,"version":false,"seq_no_primary_term":false,"explain":false,"_source":false,"fields":[{"field":"gender"},{"field":"account_number"},{"field":"firstname"},{"field":"address"},{"field":"balance"},{"field":"city"},{"field":"employer"},{"field":"state"},{"field":"age"},{"field":"email"},{"field":"lastname"}]}}}}}}, requestedTotalSize=2147483647, pageSize=null, startFrom=0)]) No newline at end of file

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

We won't have FILTER->IS NOT NULL($4) push down for common aggregate push down, after this PR: #4843

Could you check if we can remove them as well for dedup push down?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

done


// 3. Push an Aggregate
final List<RexNode> newDedupColumns = RexUtil.apply(mappingForDedupColumns, dedupColumns);
relBuilder.aggregate(relBuilder.groupKey(newDedupColumns), relBuilder.literalAgg(dedupNumer));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

LITERAL_AGG in Calcite has totally different function as we use here. It seems to be tricky to implement this in this way. Do we have alternative approach? Or at least add some comments to notice that.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I don't have other alternative approach. I think it's safe since LITERAL_AGG usually used in Project removing for Aggregate. And PPL doesn't have explicit syntax to call LITERAL_AGG. I will add comment to elaborate.

}
Integer dedupNumber = literal.getValueAs(Integer.class);
TopHitsAggregationBuilder topHitsAggregationBuilder =
AggregationBuilders.topHits(aggFieldName).from(0).fetchSource(false).size(dedupNumber);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why we set fetchSource(false) here? The default value is true?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Update: change to use fetchSource instead of fetchField due to fetchField cannot work on OS object type.

// LinkedHashMap["name" -> "A", "category" -> "X"],
// LinkedHashMap["name" -> "A", "category" -> "X"]
// ]
List<Map<String, Object>> res =

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket? I'm thinking Is it still appropriate to translate dedup x to aggregate? Normally speaking, aggregate should return only 1 row value for the metric in each bucket. While dedup x can return more than 1 rows and it now affects the API of agg metric parser.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket?

Yeah, just top_hits metric agg does. top_hits metric agg in OpenSearch can be used in any bucket agg in DSL as a sub-aggregation, it's quite different with SQL' aggregate.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@@ -1601,15 +1600,15 @@ private static void buildDedupNotNull(
.partitionBy(dedupeFields)
.orderBy(dedupeFields)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Will dedupe work without ordering by deduped fields?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

em, it should work in non-pushdown case. maybe we could remove this orderBy in window.

@yuancuyuancuNov 26, 2025

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

It looks a little strange to me because the sort keys in a window should be the same (the partition key)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. it was not introduced by this pr. let's fix them in followup PR since it will change the plan a lot.


@ToString.Exclude private final Settings settings;

@ToString.Exclude private boolean topHitsAgg = false;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

What is this field used for? It seems it's not referred anywhere

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

yes. can be deleted in current impl

Comment on lines +1820 to +1828
@Ignore("https://github.com/opensearch-project/sql/issues/4789")
public void testDedupExpr() throws IOException {
enabledOnlyWhenPushdownIsEnabled();
String expected = loadExpectedPlan("explain_dedup_expr1.yaml");
assertYamlEqualsIgnoreId(
expected,
explainQueryYaml(
"source=opensearch-sql_test_index_account | eval new_gender = lower(gender) | dedup 1"
+ " new_gender"));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Isn't such case covered by CalcitePPLDedupIT.testDedupExpr

String.format("Unsupported push-down aggregator %s", aggCall.getAggregation()));
};
}
case LITERAL_AGG -> {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

One alternative in my mind is using a self-defined class extends calcite.Aggregate. It should also be able to leverage our AggregateAnalyzer.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Checked, so far no LITERAL_AGG can be produced because PPL doesn't support literal in aggregators.
In SQL

SELECT dept_id,
COUNT(*) as emp_count,
1 as constant_val
FROM employees
GROUP BY dept_id;

Could be rewritten to

SELECT dept_id,
COUNT(*) as emp_count,
LITERAL_AGG(1) as constant_val
FROM employees
GROUP BY dept_id;

to reduce a Project upon Aggregate.

From

Project(dept_id, emp_count, constant_val=1)
Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count)
Scan(employees)

To

Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count, compute LITERAL_AGG(1) as constant_val)
Scan(employees)

But stats only support pre-defined aggregators. cc @qianheng-aws

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

I see SubQueryRemoveRule will produce LITERAL_AGG for some(we don't have) and in subquery. But I'm not sure whether it will be triggered.

And RelBuilder's literalAgg method is public, we should avoid call that method by developers then.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

@LantaoJin
LantaoJin merged commit 5ceacb6 into opensearch-project:mainNov 26, 2025
39 checks passed
@opensearch-trigger-bot

Copy link
Copy Markdown
Contributor

The backport to 2.19-dev failed:

The process '/usr/bin/git' failed with exit code 128

To backport manually, run these commands in your terminal:

# Navigate to the root of your repositorycd$(git rev-parse --show-toplevel)# Fetch latest updates from GitHub
git fetch
# Create a new working tree
git worktree add ../.worktrees/sql/backport-2.19-dev 2.19-dev
# Navigate to the new working treepushd ../.worktrees/sql/backport-2.19-dev
# Create a new branch
git switch --create backport/backport-4844-to-2.19-dev
# Cherry-pick the merged commit of this pull request and resolve the conflicts
git cherry-pick -x --mainline 1 5ceacb6945d6bacc36e037ce4da215e1cf031b56
# Push it to GitHub
git push --set-upstream origin backport/backport-4844-to-2.19-dev
# Go back to the original working treepopd# Delete the working tree
git worktree remove ../.worktrees/sql/backport-2.19-dev

Then, create a pull request where the base branch is 2.19-dev and the compare/head branch is backport/backport-4844-to-2.19-dev.

LantaoJin added a commit to LantaoJin/search-plugins-sql that referenced this pull request Nov 27, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
(cherry picked from commit 5ceacb6)
@LantaoJinLantaoJin added the backport-manually Filed a PR to backport manually. label Nov 27, 2025
asifabashar pushed a commit to asifabashar/sql that referenced this pull request Dec 10, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

backport 2.19-devbackport-failedbackport-manuallyFiled a PR to backport manually.enhancementNew feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[ENHANCEMENT] Convert dedup pushdown to composite + top_hits

3 participants

@LantaoJin@yuancu@qianheng-aws
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Convert dedup pushdown to composite + top_hits - #4844

Merged
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797
Nov 26, 2025
Merged

Convert dedup pushdown to composite + top_hits#4844
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797

Conversation

@LantaoJin

@LantaoJinLantaoJin commented Nov 21, 2025

Copy link
Copy Markdown
Member

Description

Convert dedup pushdown to composite + top_hits
Main changes:

  1. Add a rule DedupPushdownRule to convert the dedup plan pattern to Aggregate with TopHits metrics
  2. Upgrade TopHitsParser to support parsing data from _source (fetchField cannot work for OS Object type)
  3. Corresponding API change on MetricParser: Map<String, Object> parse(Aggregation aggregation); -> List<Map<String, Object>> parse(Aggregation aggregation);

Follow-ups: Support dedup on script (expression)

Related Issues

Resolves#4797

Check List

  • New functionality includes testing.
  • New functionality has been documented.
  • New functionality has javadoc added.
  • New functionality has a user manual doc added.
  • New PPL command checklist all confirmed.
  • API changes companion pull request created.
  • Commits are signed per the DCO using --signoff or -s.
  • Public documentation issue/PR created.

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.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@LantaoJinLantaoJin added the enhancement New feature or request label Nov 21, 2025
@LantaoJinLantaoJin mentioned this pull request Nov 21, 2025
8 tasks
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
return project == null
? List.of()
: aggCall.getArgList().stream().map(project.getProjects()::get).toList();
: PlanUtils.getObjectFromLiteralAgg(aggCall) != null

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why calling this method here?

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Is it for identifying LITERAL_AGG?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. LITERAL_AGG passes the numberOfDedup here

LogicalFilter(condition=[IS NOT NULL($4)])
CalciteLogicalIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]])
physical: |
CalciteEnumerableIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]], PushDownContext=[[PROJECT->[account_number, firstname, address, balance, gender, city, employer, state, age, email, lastname], FILTER->IS NOT NULL($4), AGGREGATION->rel#:LogicalAggregate.NONE.[](input=LogicalProject#,group={0},agg#0=LITERAL_AGG(1)), LIMIT->10000], OpenSearchRequestBuilder(sourceBuilder={"from":0,"size":0,"timeout":"1m","query":{"exists":{"field":"gender","boost":1.0}},"_source":{"includes":["account_number","firstname","address","balance","gender","city","employer","state","age","email","lastname"],"excludes":[]},"aggregations":{"composite_buckets":{"composite":{"size":1000,"sources":[{"gender":{"terms":{"field":"gender.keyword","missing_bucket":false,"order":"asc"}}}]},"aggregations":{"$f1":{"top_hits":{"from":0,"size":1,"version":false,"seq_no_primary_term":false,"explain":false,"_source":false,"fields":[{"field":"gender"},{"field":"account_number"},{"field":"firstname"},{"field":"address"},{"field":"balance"},{"field":"city"},{"field":"employer"},{"field":"state"},{"field":"age"},{"field":"email"},{"field":"lastname"}]}}}}}}, requestedTotalSize=2147483647, pageSize=null, startFrom=0)]) No newline at end of file

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

We won't have FILTER->IS NOT NULL($4) push down for common aggregate push down, after this PR: #4843

Could you check if we can remove them as well for dedup push down?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

done


// 3. Push an Aggregate
final List<RexNode> newDedupColumns = RexUtil.apply(mappingForDedupColumns, dedupColumns);
relBuilder.aggregate(relBuilder.groupKey(newDedupColumns), relBuilder.literalAgg(dedupNumer));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

LITERAL_AGG in Calcite has totally different function as we use here. It seems to be tricky to implement this in this way. Do we have alternative approach? Or at least add some comments to notice that.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I don't have other alternative approach. I think it's safe since LITERAL_AGG usually used in Project removing for Aggregate. And PPL doesn't have explicit syntax to call LITERAL_AGG. I will add comment to elaborate.

}
Integer dedupNumber = literal.getValueAs(Integer.class);
TopHitsAggregationBuilder topHitsAggregationBuilder =
AggregationBuilders.topHits(aggFieldName).from(0).fetchSource(false).size(dedupNumber);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why we set fetchSource(false) here? The default value is true?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Update: change to use fetchSource instead of fetchField due to fetchField cannot work on OS object type.

// LinkedHashMap["name" -> "A", "category" -> "X"],
// LinkedHashMap["name" -> "A", "category" -> "X"]
// ]
List<Map<String, Object>> res =

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket? I'm thinking Is it still appropriate to translate dedup x to aggregate? Normally speaking, aggregate should return only 1 row value for the metric in each bucket. While dedup x can return more than 1 rows and it now affects the API of agg metric parser.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket?

Yeah, just top_hits metric agg does. top_hits metric agg in OpenSearch can be used in any bucket agg in DSL as a sub-aggregation, it's quite different with SQL' aggregate.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@@ -1601,15 +1600,15 @@ private static void buildDedupNotNull(
.partitionBy(dedupeFields)
.orderBy(dedupeFields)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Will dedupe work without ordering by deduped fields?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

em, it should work in non-pushdown case. maybe we could remove this orderBy in window.

@yuancuyuancuNov 26, 2025

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

It looks a little strange to me because the sort keys in a window should be the same (the partition key)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. it was not introduced by this pr. let's fix them in followup PR since it will change the plan a lot.


@ToString.Exclude private final Settings settings;

@ToString.Exclude private boolean topHitsAgg = false;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

What is this field used for? It seems it's not referred anywhere

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

yes. can be deleted in current impl

Comment on lines +1820 to +1828
@Ignore("https://github.com/opensearch-project/sql/issues/4789")
public void testDedupExpr() throws IOException {
enabledOnlyWhenPushdownIsEnabled();
String expected = loadExpectedPlan("explain_dedup_expr1.yaml");
assertYamlEqualsIgnoreId(
expected,
explainQueryYaml(
"source=opensearch-sql_test_index_account | eval new_gender = lower(gender) | dedup 1"
+ " new_gender"));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Isn't such case covered by CalcitePPLDedupIT.testDedupExpr

String.format("Unsupported push-down aggregator %s", aggCall.getAggregation()));
};
}
case LITERAL_AGG -> {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

One alternative in my mind is using a self-defined class extends calcite.Aggregate. It should also be able to leverage our AggregateAnalyzer.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Checked, so far no LITERAL_AGG can be produced because PPL doesn't support literal in aggregators.
In SQL

SELECT dept_id,
COUNT(*) as emp_count,
1 as constant_val
FROM employees
GROUP BY dept_id;

Could be rewritten to

SELECT dept_id,
COUNT(*) as emp_count,
LITERAL_AGG(1) as constant_val
FROM employees
GROUP BY dept_id;

to reduce a Project upon Aggregate.

From

Project(dept_id, emp_count, constant_val=1)
Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count)
Scan(employees)

To

Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count, compute LITERAL_AGG(1) as constant_val)
Scan(employees)

But stats only support pre-defined aggregators. cc @qianheng-aws

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

I see SubQueryRemoveRule will produce LITERAL_AGG for some(we don't have) and in subquery. But I'm not sure whether it will be triggered.

And RelBuilder's literalAgg method is public, we should avoid call that method by developers then.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

@LantaoJin
LantaoJin merged commit 5ceacb6 into opensearch-project:mainNov 26, 2025
39 checks passed
@opensearch-trigger-bot

Copy link
Copy Markdown
Contributor

The backport to 2.19-dev failed:

The process '/usr/bin/git' failed with exit code 128

To backport manually, run these commands in your terminal:

# Navigate to the root of your repositorycd$(git rev-parse --show-toplevel)# Fetch latest updates from GitHub
git fetch
# Create a new working tree
git worktree add ../.worktrees/sql/backport-2.19-dev 2.19-dev
# Navigate to the new working treepushd ../.worktrees/sql/backport-2.19-dev
# Create a new branch
git switch --create backport/backport-4844-to-2.19-dev
# Cherry-pick the merged commit of this pull request and resolve the conflicts
git cherry-pick -x --mainline 1 5ceacb6945d6bacc36e037ce4da215e1cf031b56
# Push it to GitHub
git push --set-upstream origin backport/backport-4844-to-2.19-dev
# Go back to the original working treepopd# Delete the working tree
git worktree remove ../.worktrees/sql/backport-2.19-dev

Then, create a pull request where the base branch is 2.19-dev and the compare/head branch is backport/backport-4844-to-2.19-dev.

LantaoJin added a commit to LantaoJin/search-plugins-sql that referenced this pull request Nov 27, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
(cherry picked from commit 5ceacb6)
@LantaoJinLantaoJin added the backport-manually Filed a PR to backport manually. label Nov 27, 2025
asifabashar pushed a commit to asifabashar/sql that referenced this pull request Dec 10, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

backport 2.19-devbackport-failedbackport-manuallyFiled a PR to backport manually.enhancementNew feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[ENHANCEMENT] Convert dedup pushdown to composite + top_hits

3 participants

@LantaoJin@yuancu@qianheng-aws
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Convert dedup pushdown to composite + top_hits - #4844

Merged
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797
Nov 26, 2025
Merged

Convert dedup pushdown to composite + top_hits#4844
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797

Conversation

@LantaoJin

@LantaoJinLantaoJin commented Nov 21, 2025

Copy link
Copy Markdown
Member

Description

Convert dedup pushdown to composite + top_hits
Main changes:

  1. Add a rule DedupPushdownRule to convert the dedup plan pattern to Aggregate with TopHits metrics
  2. Upgrade TopHitsParser to support parsing data from _source (fetchField cannot work for OS Object type)
  3. Corresponding API change on MetricParser: Map<String, Object> parse(Aggregation aggregation); -> List<Map<String, Object>> parse(Aggregation aggregation);

Follow-ups: Support dedup on script (expression)

Related Issues

Resolves#4797

Check List

  • New functionality includes testing.
  • New functionality has been documented.
  • New functionality has javadoc added.
  • New functionality has a user manual doc added.
  • New PPL command checklist all confirmed.
  • API changes companion pull request created.
  • Commits are signed per the DCO using --signoff or -s.
  • Public documentation issue/PR created.

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.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@LantaoJinLantaoJin added the enhancement New feature or request label Nov 21, 2025
@LantaoJinLantaoJin mentioned this pull request Nov 21, 2025
8 tasks
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
return project == null
? List.of()
: aggCall.getArgList().stream().map(project.getProjects()::get).toList();
: PlanUtils.getObjectFromLiteralAgg(aggCall) != null

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why calling this method here?

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Is it for identifying LITERAL_AGG?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. LITERAL_AGG passes the numberOfDedup here

LogicalFilter(condition=[IS NOT NULL($4)])
CalciteLogicalIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]])
physical: |
CalciteEnumerableIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]], PushDownContext=[[PROJECT->[account_number, firstname, address, balance, gender, city, employer, state, age, email, lastname], FILTER->IS NOT NULL($4), AGGREGATION->rel#:LogicalAggregate.NONE.[](input=LogicalProject#,group={0},agg#0=LITERAL_AGG(1)), LIMIT->10000], OpenSearchRequestBuilder(sourceBuilder={"from":0,"size":0,"timeout":"1m","query":{"exists":{"field":"gender","boost":1.0}},"_source":{"includes":["account_number","firstname","address","balance","gender","city","employer","state","age","email","lastname"],"excludes":[]},"aggregations":{"composite_buckets":{"composite":{"size":1000,"sources":[{"gender":{"terms":{"field":"gender.keyword","missing_bucket":false,"order":"asc"}}}]},"aggregations":{"$f1":{"top_hits":{"from":0,"size":1,"version":false,"seq_no_primary_term":false,"explain":false,"_source":false,"fields":[{"field":"gender"},{"field":"account_number"},{"field":"firstname"},{"field":"address"},{"field":"balance"},{"field":"city"},{"field":"employer"},{"field":"state"},{"field":"age"},{"field":"email"},{"field":"lastname"}]}}}}}}, requestedTotalSize=2147483647, pageSize=null, startFrom=0)]) No newline at end of file

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

We won't have FILTER->IS NOT NULL($4) push down for common aggregate push down, after this PR: #4843

Could you check if we can remove them as well for dedup push down?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

done


// 3. Push an Aggregate
final List<RexNode> newDedupColumns = RexUtil.apply(mappingForDedupColumns, dedupColumns);
relBuilder.aggregate(relBuilder.groupKey(newDedupColumns), relBuilder.literalAgg(dedupNumer));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

LITERAL_AGG in Calcite has totally different function as we use here. It seems to be tricky to implement this in this way. Do we have alternative approach? Or at least add some comments to notice that.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I don't have other alternative approach. I think it's safe since LITERAL_AGG usually used in Project removing for Aggregate. And PPL doesn't have explicit syntax to call LITERAL_AGG. I will add comment to elaborate.

}
Integer dedupNumber = literal.getValueAs(Integer.class);
TopHitsAggregationBuilder topHitsAggregationBuilder =
AggregationBuilders.topHits(aggFieldName).from(0).fetchSource(false).size(dedupNumber);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why we set fetchSource(false) here? The default value is true?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Update: change to use fetchSource instead of fetchField due to fetchField cannot work on OS object type.

// LinkedHashMap["name" -> "A", "category" -> "X"],
// LinkedHashMap["name" -> "A", "category" -> "X"]
// ]
List<Map<String, Object>> res =

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket? I'm thinking Is it still appropriate to translate dedup x to aggregate? Normally speaking, aggregate should return only 1 row value for the metric in each bucket. While dedup x can return more than 1 rows and it now affects the API of agg metric parser.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket?

Yeah, just top_hits metric agg does. top_hits metric agg in OpenSearch can be used in any bucket agg in DSL as a sub-aggregation, it's quite different with SQL' aggregate.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@@ -1601,15 +1600,15 @@ private static void buildDedupNotNull(
.partitionBy(dedupeFields)
.orderBy(dedupeFields)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Will dedupe work without ordering by deduped fields?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

em, it should work in non-pushdown case. maybe we could remove this orderBy in window.

@yuancuyuancuNov 26, 2025

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

It looks a little strange to me because the sort keys in a window should be the same (the partition key)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. it was not introduced by this pr. let's fix them in followup PR since it will change the plan a lot.


@ToString.Exclude private final Settings settings;

@ToString.Exclude private boolean topHitsAgg = false;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

What is this field used for? It seems it's not referred anywhere

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

yes. can be deleted in current impl

Comment on lines +1820 to +1828
@Ignore("https://github.com/opensearch-project/sql/issues/4789")
public void testDedupExpr() throws IOException {
enabledOnlyWhenPushdownIsEnabled();
String expected = loadExpectedPlan("explain_dedup_expr1.yaml");
assertYamlEqualsIgnoreId(
expected,
explainQueryYaml(
"source=opensearch-sql_test_index_account | eval new_gender = lower(gender) | dedup 1"
+ " new_gender"));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Isn't such case covered by CalcitePPLDedupIT.testDedupExpr

String.format("Unsupported push-down aggregator %s", aggCall.getAggregation()));
};
}
case LITERAL_AGG -> {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

One alternative in my mind is using a self-defined class extends calcite.Aggregate. It should also be able to leverage our AggregateAnalyzer.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Checked, so far no LITERAL_AGG can be produced because PPL doesn't support literal in aggregators.
In SQL

SELECT dept_id,
COUNT(*) as emp_count,
1 as constant_val
FROM employees
GROUP BY dept_id;

Could be rewritten to

SELECT dept_id,
COUNT(*) as emp_count,
LITERAL_AGG(1) as constant_val
FROM employees
GROUP BY dept_id;

to reduce a Project upon Aggregate.

From

Project(dept_id, emp_count, constant_val=1)
Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count)
Scan(employees)

To

Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count, compute LITERAL_AGG(1) as constant_val)
Scan(employees)

But stats only support pre-defined aggregators. cc @qianheng-aws

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

I see SubQueryRemoveRule will produce LITERAL_AGG for some(we don't have) and in subquery. But I'm not sure whether it will be triggered.

And RelBuilder's literalAgg method is public, we should avoid call that method by developers then.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

@LantaoJin
LantaoJin merged commit 5ceacb6 into opensearch-project:mainNov 26, 2025
39 checks passed
@opensearch-trigger-bot

Copy link
Copy Markdown
Contributor

The backport to 2.19-dev failed:

The process '/usr/bin/git' failed with exit code 128

To backport manually, run these commands in your terminal:

# Navigate to the root of your repositorycd$(git rev-parse --show-toplevel)# Fetch latest updates from GitHub
git fetch
# Create a new working tree
git worktree add ../.worktrees/sql/backport-2.19-dev 2.19-dev
# Navigate to the new working treepushd ../.worktrees/sql/backport-2.19-dev
# Create a new branch
git switch --create backport/backport-4844-to-2.19-dev
# Cherry-pick the merged commit of this pull request and resolve the conflicts
git cherry-pick -x --mainline 1 5ceacb6945d6bacc36e037ce4da215e1cf031b56
# Push it to GitHub
git push --set-upstream origin backport/backport-4844-to-2.19-dev
# Go back to the original working treepopd# Delete the working tree
git worktree remove ../.worktrees/sql/backport-2.19-dev

Then, create a pull request where the base branch is 2.19-dev and the compare/head branch is backport/backport-4844-to-2.19-dev.

LantaoJin added a commit to LantaoJin/search-plugins-sql that referenced this pull request Nov 27, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
(cherry picked from commit 5ceacb6)
@LantaoJinLantaoJin added the backport-manually Filed a PR to backport manually. label Nov 27, 2025
asifabashar pushed a commit to asifabashar/sql that referenced this pull request Dec 10, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

backport 2.19-devbackport-failedbackport-manuallyFiled a PR to backport manually.enhancementNew feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[ENHANCEMENT] Convert dedup pushdown to composite + top_hits

3 participants

@LantaoJin@yuancu@qianheng-aws
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Convert dedup pushdown to composite + top_hits - #4844

Merged
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797
Nov 26, 2025
Merged

Convert dedup pushdown to composite + top_hits#4844
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797

Conversation

@LantaoJin

@LantaoJinLantaoJin commented Nov 21, 2025

Copy link
Copy Markdown
Member

Description

Convert dedup pushdown to composite + top_hits
Main changes:

  1. Add a rule DedupPushdownRule to convert the dedup plan pattern to Aggregate with TopHits metrics
  2. Upgrade TopHitsParser to support parsing data from _source (fetchField cannot work for OS Object type)
  3. Corresponding API change on MetricParser: Map<String, Object> parse(Aggregation aggregation); -> List<Map<String, Object>> parse(Aggregation aggregation);

Follow-ups: Support dedup on script (expression)

Related Issues

Resolves#4797

Check List

  • New functionality includes testing.
  • New functionality has been documented.
  • New functionality has javadoc added.
  • New functionality has a user manual doc added.
  • New PPL command checklist all confirmed.
  • API changes companion pull request created.
  • Commits are signed per the DCO using --signoff or -s.
  • Public documentation issue/PR created.

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.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@LantaoJinLantaoJin added the enhancement New feature or request label Nov 21, 2025
@LantaoJinLantaoJin mentioned this pull request Nov 21, 2025
8 tasks
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
return project == null
? List.of()
: aggCall.getArgList().stream().map(project.getProjects()::get).toList();
: PlanUtils.getObjectFromLiteralAgg(aggCall) != null

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why calling this method here?

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Is it for identifying LITERAL_AGG?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. LITERAL_AGG passes the numberOfDedup here

LogicalFilter(condition=[IS NOT NULL($4)])
CalciteLogicalIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]])
physical: |
CalciteEnumerableIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]], PushDownContext=[[PROJECT->[account_number, firstname, address, balance, gender, city, employer, state, age, email, lastname], FILTER->IS NOT NULL($4), AGGREGATION->rel#:LogicalAggregate.NONE.[](input=LogicalProject#,group={0},agg#0=LITERAL_AGG(1)), LIMIT->10000], OpenSearchRequestBuilder(sourceBuilder={"from":0,"size":0,"timeout":"1m","query":{"exists":{"field":"gender","boost":1.0}},"_source":{"includes":["account_number","firstname","address","balance","gender","city","employer","state","age","email","lastname"],"excludes":[]},"aggregations":{"composite_buckets":{"composite":{"size":1000,"sources":[{"gender":{"terms":{"field":"gender.keyword","missing_bucket":false,"order":"asc"}}}]},"aggregations":{"$f1":{"top_hits":{"from":0,"size":1,"version":false,"seq_no_primary_term":false,"explain":false,"_source":false,"fields":[{"field":"gender"},{"field":"account_number"},{"field":"firstname"},{"field":"address"},{"field":"balance"},{"field":"city"},{"field":"employer"},{"field":"state"},{"field":"age"},{"field":"email"},{"field":"lastname"}]}}}}}}, requestedTotalSize=2147483647, pageSize=null, startFrom=0)]) No newline at end of file

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

We won't have FILTER->IS NOT NULL($4) push down for common aggregate push down, after this PR: #4843

Could you check if we can remove them as well for dedup push down?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

done


// 3. Push an Aggregate
final List<RexNode> newDedupColumns = RexUtil.apply(mappingForDedupColumns, dedupColumns);
relBuilder.aggregate(relBuilder.groupKey(newDedupColumns), relBuilder.literalAgg(dedupNumer));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

LITERAL_AGG in Calcite has totally different function as we use here. It seems to be tricky to implement this in this way. Do we have alternative approach? Or at least add some comments to notice that.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I don't have other alternative approach. I think it's safe since LITERAL_AGG usually used in Project removing for Aggregate. And PPL doesn't have explicit syntax to call LITERAL_AGG. I will add comment to elaborate.

}
Integer dedupNumber = literal.getValueAs(Integer.class);
TopHitsAggregationBuilder topHitsAggregationBuilder =
AggregationBuilders.topHits(aggFieldName).from(0).fetchSource(false).size(dedupNumber);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why we set fetchSource(false) here? The default value is true?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Update: change to use fetchSource instead of fetchField due to fetchField cannot work on OS object type.

// LinkedHashMap["name" -> "A", "category" -> "X"],
// LinkedHashMap["name" -> "A", "category" -> "X"]
// ]
List<Map<String, Object>> res =

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket? I'm thinking Is it still appropriate to translate dedup x to aggregate? Normally speaking, aggregate should return only 1 row value for the metric in each bucket. While dedup x can return more than 1 rows and it now affects the API of agg metric parser.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket?

Yeah, just top_hits metric agg does. top_hits metric agg in OpenSearch can be used in any bucket agg in DSL as a sub-aggregation, it's quite different with SQL' aggregate.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@@ -1601,15 +1600,15 @@ private static void buildDedupNotNull(
.partitionBy(dedupeFields)
.orderBy(dedupeFields)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Will dedupe work without ordering by deduped fields?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

em, it should work in non-pushdown case. maybe we could remove this orderBy in window.

@yuancuyuancuNov 26, 2025

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

It looks a little strange to me because the sort keys in a window should be the same (the partition key)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. it was not introduced by this pr. let's fix them in followup PR since it will change the plan a lot.


@ToString.Exclude private final Settings settings;

@ToString.Exclude private boolean topHitsAgg = false;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

What is this field used for? It seems it's not referred anywhere

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

yes. can be deleted in current impl

Comment on lines +1820 to +1828
@Ignore("https://github.com/opensearch-project/sql/issues/4789")
public void testDedupExpr() throws IOException {
enabledOnlyWhenPushdownIsEnabled();
String expected = loadExpectedPlan("explain_dedup_expr1.yaml");
assertYamlEqualsIgnoreId(
expected,
explainQueryYaml(
"source=opensearch-sql_test_index_account | eval new_gender = lower(gender) | dedup 1"
+ " new_gender"));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Isn't such case covered by CalcitePPLDedupIT.testDedupExpr

String.format("Unsupported push-down aggregator %s", aggCall.getAggregation()));
};
}
case LITERAL_AGG -> {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

One alternative in my mind is using a self-defined class extends calcite.Aggregate. It should also be able to leverage our AggregateAnalyzer.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Checked, so far no LITERAL_AGG can be produced because PPL doesn't support literal in aggregators.
In SQL

SELECT dept_id,
COUNT(*) as emp_count,
1 as constant_val
FROM employees
GROUP BY dept_id;

Could be rewritten to

SELECT dept_id,
COUNT(*) as emp_count,
LITERAL_AGG(1) as constant_val
FROM employees
GROUP BY dept_id;

to reduce a Project upon Aggregate.

From

Project(dept_id, emp_count, constant_val=1)
Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count)
Scan(employees)

To

Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count, compute LITERAL_AGG(1) as constant_val)
Scan(employees)

But stats only support pre-defined aggregators. cc @qianheng-aws

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

I see SubQueryRemoveRule will produce LITERAL_AGG for some(we don't have) and in subquery. But I'm not sure whether it will be triggered.

And RelBuilder's literalAgg method is public, we should avoid call that method by developers then.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

@LantaoJin
LantaoJin merged commit 5ceacb6 into opensearch-project:mainNov 26, 2025
39 checks passed
@opensearch-trigger-bot

Copy link
Copy Markdown
Contributor

The backport to 2.19-dev failed:

The process '/usr/bin/git' failed with exit code 128

To backport manually, run these commands in your terminal:

# Navigate to the root of your repositorycd$(git rev-parse --show-toplevel)# Fetch latest updates from GitHub
git fetch
# Create a new working tree
git worktree add ../.worktrees/sql/backport-2.19-dev 2.19-dev
# Navigate to the new working treepushd ../.worktrees/sql/backport-2.19-dev
# Create a new branch
git switch --create backport/backport-4844-to-2.19-dev
# Cherry-pick the merged commit of this pull request and resolve the conflicts
git cherry-pick -x --mainline 1 5ceacb6945d6bacc36e037ce4da215e1cf031b56
# Push it to GitHub
git push --set-upstream origin backport/backport-4844-to-2.19-dev
# Go back to the original working treepopd# Delete the working tree
git worktree remove ../.worktrees/sql/backport-2.19-dev

Then, create a pull request where the base branch is 2.19-dev and the compare/head branch is backport/backport-4844-to-2.19-dev.

LantaoJin added a commit to LantaoJin/search-plugins-sql that referenced this pull request Nov 27, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
(cherry picked from commit 5ceacb6)
@LantaoJinLantaoJin added the backport-manually Filed a PR to backport manually. label Nov 27, 2025
asifabashar pushed a commit to asifabashar/sql that referenced this pull request Dec 10, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

backport 2.19-devbackport-failedbackport-manuallyFiled a PR to backport manually.enhancementNew feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[ENHANCEMENT] Convert dedup pushdown to composite + top_hits

3 participants

@LantaoJin@yuancu@qianheng-aws
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Convert dedup pushdown to composite + top_hits - #4844

Merged
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797
Nov 26, 2025
Merged

Convert dedup pushdown to composite + top_hits#4844
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797

Conversation

@LantaoJin

@LantaoJinLantaoJin commented Nov 21, 2025

Copy link
Copy Markdown
Member

Description

Convert dedup pushdown to composite + top_hits
Main changes:

  1. Add a rule DedupPushdownRule to convert the dedup plan pattern to Aggregate with TopHits metrics
  2. Upgrade TopHitsParser to support parsing data from _source (fetchField cannot work for OS Object type)
  3. Corresponding API change on MetricParser: Map<String, Object> parse(Aggregation aggregation); -> List<Map<String, Object>> parse(Aggregation aggregation);

Follow-ups: Support dedup on script (expression)

Related Issues

Resolves#4797

Check List

  • New functionality includes testing.
  • New functionality has been documented.
  • New functionality has javadoc added.
  • New functionality has a user manual doc added.
  • New PPL command checklist all confirmed.
  • API changes companion pull request created.
  • Commits are signed per the DCO using --signoff or -s.
  • Public documentation issue/PR created.

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.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@LantaoJinLantaoJin added the enhancement New feature or request label Nov 21, 2025
@LantaoJinLantaoJin mentioned this pull request Nov 21, 2025
8 tasks
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
return project == null
? List.of()
: aggCall.getArgList().stream().map(project.getProjects()::get).toList();
: PlanUtils.getObjectFromLiteralAgg(aggCall) != null

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why calling this method here?

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Is it for identifying LITERAL_AGG?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. LITERAL_AGG passes the numberOfDedup here

LogicalFilter(condition=[IS NOT NULL($4)])
CalciteLogicalIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]])
physical: |
CalciteEnumerableIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]], PushDownContext=[[PROJECT->[account_number, firstname, address, balance, gender, city, employer, state, age, email, lastname], FILTER->IS NOT NULL($4), AGGREGATION->rel#:LogicalAggregate.NONE.[](input=LogicalProject#,group={0},agg#0=LITERAL_AGG(1)), LIMIT->10000], OpenSearchRequestBuilder(sourceBuilder={"from":0,"size":0,"timeout":"1m","query":{"exists":{"field":"gender","boost":1.0}},"_source":{"includes":["account_number","firstname","address","balance","gender","city","employer","state","age","email","lastname"],"excludes":[]},"aggregations":{"composite_buckets":{"composite":{"size":1000,"sources":[{"gender":{"terms":{"field":"gender.keyword","missing_bucket":false,"order":"asc"}}}]},"aggregations":{"$f1":{"top_hits":{"from":0,"size":1,"version":false,"seq_no_primary_term":false,"explain":false,"_source":false,"fields":[{"field":"gender"},{"field":"account_number"},{"field":"firstname"},{"field":"address"},{"field":"balance"},{"field":"city"},{"field":"employer"},{"field":"state"},{"field":"age"},{"field":"email"},{"field":"lastname"}]}}}}}}, requestedTotalSize=2147483647, pageSize=null, startFrom=0)]) No newline at end of file

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

We won't have FILTER->IS NOT NULL($4) push down for common aggregate push down, after this PR: #4843

Could you check if we can remove them as well for dedup push down?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

done


// 3. Push an Aggregate
final List<RexNode> newDedupColumns = RexUtil.apply(mappingForDedupColumns, dedupColumns);
relBuilder.aggregate(relBuilder.groupKey(newDedupColumns), relBuilder.literalAgg(dedupNumer));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

LITERAL_AGG in Calcite has totally different function as we use here. It seems to be tricky to implement this in this way. Do we have alternative approach? Or at least add some comments to notice that.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I don't have other alternative approach. I think it's safe since LITERAL_AGG usually used in Project removing for Aggregate. And PPL doesn't have explicit syntax to call LITERAL_AGG. I will add comment to elaborate.

}
Integer dedupNumber = literal.getValueAs(Integer.class);
TopHitsAggregationBuilder topHitsAggregationBuilder =
AggregationBuilders.topHits(aggFieldName).from(0).fetchSource(false).size(dedupNumber);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why we set fetchSource(false) here? The default value is true?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Update: change to use fetchSource instead of fetchField due to fetchField cannot work on OS object type.

// LinkedHashMap["name" -> "A", "category" -> "X"],
// LinkedHashMap["name" -> "A", "category" -> "X"]
// ]
List<Map<String, Object>> res =

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket? I'm thinking Is it still appropriate to translate dedup x to aggregate? Normally speaking, aggregate should return only 1 row value for the metric in each bucket. While dedup x can return more than 1 rows and it now affects the API of agg metric parser.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket?

Yeah, just top_hits metric agg does. top_hits metric agg in OpenSearch can be used in any bucket agg in DSL as a sub-aggregation, it's quite different with SQL' aggregate.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@@ -1601,15 +1600,15 @@ private static void buildDedupNotNull(
.partitionBy(dedupeFields)
.orderBy(dedupeFields)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Will dedupe work without ordering by deduped fields?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

em, it should work in non-pushdown case. maybe we could remove this orderBy in window.

@yuancuyuancuNov 26, 2025

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

It looks a little strange to me because the sort keys in a window should be the same (the partition key)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. it was not introduced by this pr. let's fix them in followup PR since it will change the plan a lot.


@ToString.Exclude private final Settings settings;

@ToString.Exclude private boolean topHitsAgg = false;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

What is this field used for? It seems it's not referred anywhere

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

yes. can be deleted in current impl

Comment on lines +1820 to +1828
@Ignore("https://github.com/opensearch-project/sql/issues/4789")
public void testDedupExpr() throws IOException {
enabledOnlyWhenPushdownIsEnabled();
String expected = loadExpectedPlan("explain_dedup_expr1.yaml");
assertYamlEqualsIgnoreId(
expected,
explainQueryYaml(
"source=opensearch-sql_test_index_account | eval new_gender = lower(gender) | dedup 1"
+ " new_gender"));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Isn't such case covered by CalcitePPLDedupIT.testDedupExpr

String.format("Unsupported push-down aggregator %s", aggCall.getAggregation()));
};
}
case LITERAL_AGG -> {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

One alternative in my mind is using a self-defined class extends calcite.Aggregate. It should also be able to leverage our AggregateAnalyzer.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Checked, so far no LITERAL_AGG can be produced because PPL doesn't support literal in aggregators.
In SQL

SELECT dept_id,
COUNT(*) as emp_count,
1 as constant_val
FROM employees
GROUP BY dept_id;

Could be rewritten to

SELECT dept_id,
COUNT(*) as emp_count,
LITERAL_AGG(1) as constant_val
FROM employees
GROUP BY dept_id;

to reduce a Project upon Aggregate.

From

Project(dept_id, emp_count, constant_val=1)
Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count)
Scan(employees)

To

Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count, compute LITERAL_AGG(1) as constant_val)
Scan(employees)

But stats only support pre-defined aggregators. cc @qianheng-aws

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

I see SubQueryRemoveRule will produce LITERAL_AGG for some(we don't have) and in subquery. But I'm not sure whether it will be triggered.

And RelBuilder's literalAgg method is public, we should avoid call that method by developers then.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

@LantaoJin
LantaoJin merged commit 5ceacb6 into opensearch-project:mainNov 26, 2025
39 checks passed
@opensearch-trigger-bot

Copy link
Copy Markdown
Contributor

The backport to 2.19-dev failed:

The process '/usr/bin/git' failed with exit code 128

To backport manually, run these commands in your terminal:

# Navigate to the root of your repositorycd$(git rev-parse --show-toplevel)# Fetch latest updates from GitHub
git fetch
# Create a new working tree
git worktree add ../.worktrees/sql/backport-2.19-dev 2.19-dev
# Navigate to the new working treepushd ../.worktrees/sql/backport-2.19-dev
# Create a new branch
git switch --create backport/backport-4844-to-2.19-dev
# Cherry-pick the merged commit of this pull request and resolve the conflicts
git cherry-pick -x --mainline 1 5ceacb6945d6bacc36e037ce4da215e1cf031b56
# Push it to GitHub
git push --set-upstream origin backport/backport-4844-to-2.19-dev
# Go back to the original working treepopd# Delete the working tree
git worktree remove ../.worktrees/sql/backport-2.19-dev

Then, create a pull request where the base branch is 2.19-dev and the compare/head branch is backport/backport-4844-to-2.19-dev.

LantaoJin added a commit to LantaoJin/search-plugins-sql that referenced this pull request Nov 27, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
(cherry picked from commit 5ceacb6)
@LantaoJinLantaoJin added the backport-manually Filed a PR to backport manually. label Nov 27, 2025
asifabashar pushed a commit to asifabashar/sql that referenced this pull request Dec 10, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

backport 2.19-devbackport-failedbackport-manuallyFiled a PR to backport manually.enhancementNew feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[ENHANCEMENT] Convert dedup pushdown to composite + top_hits

3 participants

@LantaoJin@yuancu@qianheng-aws
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Convert dedup pushdown to composite + top_hits - #4844

Merged
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797
Nov 26, 2025
Merged

Convert dedup pushdown to composite + top_hits#4844
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797

Conversation

@LantaoJin

@LantaoJinLantaoJin commented Nov 21, 2025

Copy link
Copy Markdown
Member

Description

Convert dedup pushdown to composite + top_hits
Main changes:

  1. Add a rule DedupPushdownRule to convert the dedup plan pattern to Aggregate with TopHits metrics
  2. Upgrade TopHitsParser to support parsing data from _source (fetchField cannot work for OS Object type)
  3. Corresponding API change on MetricParser: Map<String, Object> parse(Aggregation aggregation); -> List<Map<String, Object>> parse(Aggregation aggregation);

Follow-ups: Support dedup on script (expression)

Related Issues

Resolves#4797

Check List

  • New functionality includes testing.
  • New functionality has been documented.
  • New functionality has javadoc added.
  • New functionality has a user manual doc added.
  • New PPL command checklist all confirmed.
  • API changes companion pull request created.
  • Commits are signed per the DCO using --signoff or -s.
  • Public documentation issue/PR created.

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.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@LantaoJinLantaoJin added the enhancement New feature or request label Nov 21, 2025
@LantaoJinLantaoJin mentioned this pull request Nov 21, 2025
8 tasks
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
return project == null
? List.of()
: aggCall.getArgList().stream().map(project.getProjects()::get).toList();
: PlanUtils.getObjectFromLiteralAgg(aggCall) != null

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why calling this method here?

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Is it for identifying LITERAL_AGG?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. LITERAL_AGG passes the numberOfDedup here

LogicalFilter(condition=[IS NOT NULL($4)])
CalciteLogicalIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]])
physical: |
CalciteEnumerableIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]], PushDownContext=[[PROJECT->[account_number, firstname, address, balance, gender, city, employer, state, age, email, lastname], FILTER->IS NOT NULL($4), AGGREGATION->rel#:LogicalAggregate.NONE.[](input=LogicalProject#,group={0},agg#0=LITERAL_AGG(1)), LIMIT->10000], OpenSearchRequestBuilder(sourceBuilder={"from":0,"size":0,"timeout":"1m","query":{"exists":{"field":"gender","boost":1.0}},"_source":{"includes":["account_number","firstname","address","balance","gender","city","employer","state","age","email","lastname"],"excludes":[]},"aggregations":{"composite_buckets":{"composite":{"size":1000,"sources":[{"gender":{"terms":{"field":"gender.keyword","missing_bucket":false,"order":"asc"}}}]},"aggregations":{"$f1":{"top_hits":{"from":0,"size":1,"version":false,"seq_no_primary_term":false,"explain":false,"_source":false,"fields":[{"field":"gender"},{"field":"account_number"},{"field":"firstname"},{"field":"address"},{"field":"balance"},{"field":"city"},{"field":"employer"},{"field":"state"},{"field":"age"},{"field":"email"},{"field":"lastname"}]}}}}}}, requestedTotalSize=2147483647, pageSize=null, startFrom=0)]) No newline at end of file

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

We won't have FILTER->IS NOT NULL($4) push down for common aggregate push down, after this PR: #4843

Could you check if we can remove them as well for dedup push down?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

done


// 3. Push an Aggregate
final List<RexNode> newDedupColumns = RexUtil.apply(mappingForDedupColumns, dedupColumns);
relBuilder.aggregate(relBuilder.groupKey(newDedupColumns), relBuilder.literalAgg(dedupNumer));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

LITERAL_AGG in Calcite has totally different function as we use here. It seems to be tricky to implement this in this way. Do we have alternative approach? Or at least add some comments to notice that.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I don't have other alternative approach. I think it's safe since LITERAL_AGG usually used in Project removing for Aggregate. And PPL doesn't have explicit syntax to call LITERAL_AGG. I will add comment to elaborate.

}
Integer dedupNumber = literal.getValueAs(Integer.class);
TopHitsAggregationBuilder topHitsAggregationBuilder =
AggregationBuilders.topHits(aggFieldName).from(0).fetchSource(false).size(dedupNumber);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why we set fetchSource(false) here? The default value is true?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Update: change to use fetchSource instead of fetchField due to fetchField cannot work on OS object type.

// LinkedHashMap["name" -> "A", "category" -> "X"],
// LinkedHashMap["name" -> "A", "category" -> "X"]
// ]
List<Map<String, Object>> res =

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket? I'm thinking Is it still appropriate to translate dedup x to aggregate? Normally speaking, aggregate should return only 1 row value for the metric in each bucket. While dedup x can return more than 1 rows and it now affects the API of agg metric parser.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket?

Yeah, just top_hits metric agg does. top_hits metric agg in OpenSearch can be used in any bucket agg in DSL as a sub-aggregation, it's quite different with SQL' aggregate.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@@ -1601,15 +1600,15 @@ private static void buildDedupNotNull(
.partitionBy(dedupeFields)
.orderBy(dedupeFields)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Will dedupe work without ordering by deduped fields?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

em, it should work in non-pushdown case. maybe we could remove this orderBy in window.

@yuancuyuancuNov 26, 2025

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

It looks a little strange to me because the sort keys in a window should be the same (the partition key)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. it was not introduced by this pr. let's fix them in followup PR since it will change the plan a lot.


@ToString.Exclude private final Settings settings;

@ToString.Exclude private boolean topHitsAgg = false;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

What is this field used for? It seems it's not referred anywhere

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

yes. can be deleted in current impl

Comment on lines +1820 to +1828
@Ignore("https://github.com/opensearch-project/sql/issues/4789")
public void testDedupExpr() throws IOException {
enabledOnlyWhenPushdownIsEnabled();
String expected = loadExpectedPlan("explain_dedup_expr1.yaml");
assertYamlEqualsIgnoreId(
expected,
explainQueryYaml(
"source=opensearch-sql_test_index_account | eval new_gender = lower(gender) | dedup 1"
+ " new_gender"));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Isn't such case covered by CalcitePPLDedupIT.testDedupExpr

String.format("Unsupported push-down aggregator %s", aggCall.getAggregation()));
};
}
case LITERAL_AGG -> {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

One alternative in my mind is using a self-defined class extends calcite.Aggregate. It should also be able to leverage our AggregateAnalyzer.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Checked, so far no LITERAL_AGG can be produced because PPL doesn't support literal in aggregators.
In SQL

SELECT dept_id,
COUNT(*) as emp_count,
1 as constant_val
FROM employees
GROUP BY dept_id;

Could be rewritten to

SELECT dept_id,
COUNT(*) as emp_count,
LITERAL_AGG(1) as constant_val
FROM employees
GROUP BY dept_id;

to reduce a Project upon Aggregate.

From

Project(dept_id, emp_count, constant_val=1)
Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count)
Scan(employees)

To

Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count, compute LITERAL_AGG(1) as constant_val)
Scan(employees)

But stats only support pre-defined aggregators. cc @qianheng-aws

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

I see SubQueryRemoveRule will produce LITERAL_AGG for some(we don't have) and in subquery. But I'm not sure whether it will be triggered.

And RelBuilder's literalAgg method is public, we should avoid call that method by developers then.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

@LantaoJin
LantaoJin merged commit 5ceacb6 into opensearch-project:mainNov 26, 2025
39 checks passed
@opensearch-trigger-bot

Copy link
Copy Markdown
Contributor

The backport to 2.19-dev failed:

The process '/usr/bin/git' failed with exit code 128

To backport manually, run these commands in your terminal:

# Navigate to the root of your repositorycd$(git rev-parse --show-toplevel)# Fetch latest updates from GitHub
git fetch
# Create a new working tree
git worktree add ../.worktrees/sql/backport-2.19-dev 2.19-dev
# Navigate to the new working treepushd ../.worktrees/sql/backport-2.19-dev
# Create a new branch
git switch --create backport/backport-4844-to-2.19-dev
# Cherry-pick the merged commit of this pull request and resolve the conflicts
git cherry-pick -x --mainline 1 5ceacb6945d6bacc36e037ce4da215e1cf031b56
# Push it to GitHub
git push --set-upstream origin backport/backport-4844-to-2.19-dev
# Go back to the original working treepopd# Delete the working tree
git worktree remove ../.worktrees/sql/backport-2.19-dev

Then, create a pull request where the base branch is 2.19-dev and the compare/head branch is backport/backport-4844-to-2.19-dev.

LantaoJin added a commit to LantaoJin/search-plugins-sql that referenced this pull request Nov 27, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
(cherry picked from commit 5ceacb6)
@LantaoJinLantaoJin added the backport-manually Filed a PR to backport manually. label Nov 27, 2025
asifabashar pushed a commit to asifabashar/sql that referenced this pull request Dec 10, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

backport 2.19-devbackport-failedbackport-manuallyFiled a PR to backport manually.enhancementNew feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[ENHANCEMENT] Convert dedup pushdown to composite + top_hits

3 participants

@LantaoJin@yuancu@qianheng-aws
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Convert dedup pushdown to composite + top_hits - #4844

Merged
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797
Nov 26, 2025
Merged

Convert dedup pushdown to composite + top_hits#4844
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797

Conversation

@LantaoJin

@LantaoJinLantaoJin commented Nov 21, 2025

Copy link
Copy Markdown
Member

Description

Convert dedup pushdown to composite + top_hits
Main changes:

  1. Add a rule DedupPushdownRule to convert the dedup plan pattern to Aggregate with TopHits metrics
  2. Upgrade TopHitsParser to support parsing data from _source (fetchField cannot work for OS Object type)
  3. Corresponding API change on MetricParser: Map<String, Object> parse(Aggregation aggregation); -> List<Map<String, Object>> parse(Aggregation aggregation);

Follow-ups: Support dedup on script (expression)

Related Issues

Resolves#4797

Check List

  • New functionality includes testing.
  • New functionality has been documented.
  • New functionality has javadoc added.
  • New functionality has a user manual doc added.
  • New PPL command checklist all confirmed.
  • API changes companion pull request created.
  • Commits are signed per the DCO using --signoff or -s.
  • Public documentation issue/PR created.

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.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@LantaoJinLantaoJin added the enhancement New feature or request label Nov 21, 2025
@LantaoJinLantaoJin mentioned this pull request Nov 21, 2025
8 tasks
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
return project == null
? List.of()
: aggCall.getArgList().stream().map(project.getProjects()::get).toList();
: PlanUtils.getObjectFromLiteralAgg(aggCall) != null

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why calling this method here?

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Is it for identifying LITERAL_AGG?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. LITERAL_AGG passes the numberOfDedup here

LogicalFilter(condition=[IS NOT NULL($4)])
CalciteLogicalIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]])
physical: |
CalciteEnumerableIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]], PushDownContext=[[PROJECT->[account_number, firstname, address, balance, gender, city, employer, state, age, email, lastname], FILTER->IS NOT NULL($4), AGGREGATION->rel#:LogicalAggregate.NONE.[](input=LogicalProject#,group={0},agg#0=LITERAL_AGG(1)), LIMIT->10000], OpenSearchRequestBuilder(sourceBuilder={"from":0,"size":0,"timeout":"1m","query":{"exists":{"field":"gender","boost":1.0}},"_source":{"includes":["account_number","firstname","address","balance","gender","city","employer","state","age","email","lastname"],"excludes":[]},"aggregations":{"composite_buckets":{"composite":{"size":1000,"sources":[{"gender":{"terms":{"field":"gender.keyword","missing_bucket":false,"order":"asc"}}}]},"aggregations":{"$f1":{"top_hits":{"from":0,"size":1,"version":false,"seq_no_primary_term":false,"explain":false,"_source":false,"fields":[{"field":"gender"},{"field":"account_number"},{"field":"firstname"},{"field":"address"},{"field":"balance"},{"field":"city"},{"field":"employer"},{"field":"state"},{"field":"age"},{"field":"email"},{"field":"lastname"}]}}}}}}, requestedTotalSize=2147483647, pageSize=null, startFrom=0)]) No newline at end of file

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

We won't have FILTER->IS NOT NULL($4) push down for common aggregate push down, after this PR: #4843

Could you check if we can remove them as well for dedup push down?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

done


// 3. Push an Aggregate
final List<RexNode> newDedupColumns = RexUtil.apply(mappingForDedupColumns, dedupColumns);
relBuilder.aggregate(relBuilder.groupKey(newDedupColumns), relBuilder.literalAgg(dedupNumer));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

LITERAL_AGG in Calcite has totally different function as we use here. It seems to be tricky to implement this in this way. Do we have alternative approach? Or at least add some comments to notice that.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I don't have other alternative approach. I think it's safe since LITERAL_AGG usually used in Project removing for Aggregate. And PPL doesn't have explicit syntax to call LITERAL_AGG. I will add comment to elaborate.

}
Integer dedupNumber = literal.getValueAs(Integer.class);
TopHitsAggregationBuilder topHitsAggregationBuilder =
AggregationBuilders.topHits(aggFieldName).from(0).fetchSource(false).size(dedupNumber);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why we set fetchSource(false) here? The default value is true?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Update: change to use fetchSource instead of fetchField due to fetchField cannot work on OS object type.

// LinkedHashMap["name" -> "A", "category" -> "X"],
// LinkedHashMap["name" -> "A", "category" -> "X"]
// ]
List<Map<String, Object>> res =

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket? I'm thinking Is it still appropriate to translate dedup x to aggregate? Normally speaking, aggregate should return only 1 row value for the metric in each bucket. While dedup x can return more than 1 rows and it now affects the API of agg metric parser.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket?

Yeah, just top_hits metric agg does. top_hits metric agg in OpenSearch can be used in any bucket agg in DSL as a sub-aggregation, it's quite different with SQL' aggregate.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@@ -1601,15 +1600,15 @@ private static void buildDedupNotNull(
.partitionBy(dedupeFields)
.orderBy(dedupeFields)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Will dedupe work without ordering by deduped fields?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

em, it should work in non-pushdown case. maybe we could remove this orderBy in window.

@yuancuyuancuNov 26, 2025

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

It looks a little strange to me because the sort keys in a window should be the same (the partition key)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. it was not introduced by this pr. let's fix them in followup PR since it will change the plan a lot.


@ToString.Exclude private final Settings settings;

@ToString.Exclude private boolean topHitsAgg = false;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

What is this field used for? It seems it's not referred anywhere

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

yes. can be deleted in current impl

Comment on lines +1820 to +1828
@Ignore("https://github.com/opensearch-project/sql/issues/4789")
public void testDedupExpr() throws IOException {
enabledOnlyWhenPushdownIsEnabled();
String expected = loadExpectedPlan("explain_dedup_expr1.yaml");
assertYamlEqualsIgnoreId(
expected,
explainQueryYaml(
"source=opensearch-sql_test_index_account | eval new_gender = lower(gender) | dedup 1"
+ " new_gender"));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Isn't such case covered by CalcitePPLDedupIT.testDedupExpr

String.format("Unsupported push-down aggregator %s", aggCall.getAggregation()));
};
}
case LITERAL_AGG -> {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

One alternative in my mind is using a self-defined class extends calcite.Aggregate. It should also be able to leverage our AggregateAnalyzer.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Checked, so far no LITERAL_AGG can be produced because PPL doesn't support literal in aggregators.
In SQL

SELECT dept_id,
COUNT(*) as emp_count,
1 as constant_val
FROM employees
GROUP BY dept_id;

Could be rewritten to

SELECT dept_id,
COUNT(*) as emp_count,
LITERAL_AGG(1) as constant_val
FROM employees
GROUP BY dept_id;

to reduce a Project upon Aggregate.

From

Project(dept_id, emp_count, constant_val=1)
Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count)
Scan(employees)

To

Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count, compute LITERAL_AGG(1) as constant_val)
Scan(employees)

But stats only support pre-defined aggregators. cc @qianheng-aws

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

I see SubQueryRemoveRule will produce LITERAL_AGG for some(we don't have) and in subquery. But I'm not sure whether it will be triggered.

And RelBuilder's literalAgg method is public, we should avoid call that method by developers then.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

@LantaoJin
LantaoJin merged commit 5ceacb6 into opensearch-project:mainNov 26, 2025
39 checks passed
@opensearch-trigger-bot

Copy link
Copy Markdown
Contributor

The backport to 2.19-dev failed:

The process '/usr/bin/git' failed with exit code 128

To backport manually, run these commands in your terminal:

# Navigate to the root of your repositorycd$(git rev-parse --show-toplevel)# Fetch latest updates from GitHub
git fetch
# Create a new working tree
git worktree add ../.worktrees/sql/backport-2.19-dev 2.19-dev
# Navigate to the new working treepushd ../.worktrees/sql/backport-2.19-dev
# Create a new branch
git switch --create backport/backport-4844-to-2.19-dev
# Cherry-pick the merged commit of this pull request and resolve the conflicts
git cherry-pick -x --mainline 1 5ceacb6945d6bacc36e037ce4da215e1cf031b56
# Push it to GitHub
git push --set-upstream origin backport/backport-4844-to-2.19-dev
# Go back to the original working treepopd# Delete the working tree
git worktree remove ../.worktrees/sql/backport-2.19-dev

Then, create a pull request where the base branch is 2.19-dev and the compare/head branch is backport/backport-4844-to-2.19-dev.

LantaoJin added a commit to LantaoJin/search-plugins-sql that referenced this pull request Nov 27, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
(cherry picked from commit 5ceacb6)
@LantaoJinLantaoJin added the backport-manually Filed a PR to backport manually. label Nov 27, 2025
asifabashar pushed a commit to asifabashar/sql that referenced this pull request Dec 10, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

backport 2.19-devbackport-failedbackport-manuallyFiled a PR to backport manually.enhancementNew feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[ENHANCEMENT] Convert dedup pushdown to composite + top_hits

3 participants

@LantaoJin@yuancu@qianheng-aws
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Convert dedup pushdown to composite + top_hits - #4844

Merged
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797
Nov 26, 2025
Merged

Convert dedup pushdown to composite + top_hits#4844
LantaoJin merged 12 commits into
opensearch-project:mainfrom
LantaoJin:pr/issues/4797

Conversation

@LantaoJin

@LantaoJinLantaoJin commented Nov 21, 2025

Copy link
Copy Markdown
Member

Description

Convert dedup pushdown to composite + top_hits
Main changes:

  1. Add a rule DedupPushdownRule to convert the dedup plan pattern to Aggregate with TopHits metrics
  2. Upgrade TopHitsParser to support parsing data from _source (fetchField cannot work for OS Object type)
  3. Corresponding API change on MetricParser: Map<String, Object> parse(Aggregation aggregation); -> List<Map<String, Object>> parse(Aggregation aggregation);

Follow-ups: Support dedup on script (expression)

Related Issues

Resolves#4797

Check List

  • New functionality includes testing.
  • New functionality has been documented.
  • New functionality has javadoc added.
  • New functionality has a user manual doc added.
  • New PPL command checklist all confirmed.
  • API changes companion pull request created.
  • Commits are signed per the DCO using --signoff or -s.
  • Public documentation issue/PR created.

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.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@LantaoJinLantaoJin added the enhancement New feature or request label Nov 21, 2025
@LantaoJinLantaoJin mentioned this pull request Nov 21, 2025
8 tasks
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
return project == null
? List.of()
: aggCall.getArgList().stream().map(project.getProjects()::get).toList();
: PlanUtils.getObjectFromLiteralAgg(aggCall) != null

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why calling this method here?

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Is it for identifying LITERAL_AGG?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. LITERAL_AGG passes the numberOfDedup here

LogicalFilter(condition=[IS NOT NULL($4)])
CalciteLogicalIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]])
physical: |
CalciteEnumerableIndexScan(table=[[OpenSearch, opensearch-sql_test_index_account]], PushDownContext=[[PROJECT->[account_number, firstname, address, balance, gender, city, employer, state, age, email, lastname], FILTER->IS NOT NULL($4), AGGREGATION->rel#:LogicalAggregate.NONE.[](input=LogicalProject#,group={0},agg#0=LITERAL_AGG(1)), LIMIT->10000], OpenSearchRequestBuilder(sourceBuilder={"from":0,"size":0,"timeout":"1m","query":{"exists":{"field":"gender","boost":1.0}},"_source":{"includes":["account_number","firstname","address","balance","gender","city","employer","state","age","email","lastname"],"excludes":[]},"aggregations":{"composite_buckets":{"composite":{"size":1000,"sources":[{"gender":{"terms":{"field":"gender.keyword","missing_bucket":false,"order":"asc"}}}]},"aggregations":{"$f1":{"top_hits":{"from":0,"size":1,"version":false,"seq_no_primary_term":false,"explain":false,"_source":false,"fields":[{"field":"gender"},{"field":"account_number"},{"field":"firstname"},{"field":"address"},{"field":"balance"},{"field":"city"},{"field":"employer"},{"field":"state"},{"field":"age"},{"field":"email"},{"field":"lastname"}]}}}}}}, requestedTotalSize=2147483647, pageSize=null, startFrom=0)]) No newline at end of file

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

We won't have FILTER->IS NOT NULL($4) push down for common aggregate push down, after this PR: #4843

Could you check if we can remove them as well for dedup push down?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

done


// 3. Push an Aggregate
final List<RexNode> newDedupColumns = RexUtil.apply(mappingForDedupColumns, dedupColumns);
relBuilder.aggregate(relBuilder.groupKey(newDedupColumns), relBuilder.literalAgg(dedupNumer));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

LITERAL_AGG in Calcite has totally different function as we use here. It seems to be tricky to implement this in this way. Do we have alternative approach? Or at least add some comments to notice that.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

I don't have other alternative approach. I think it's safe since LITERAL_AGG usually used in Project removing for Aggregate. And PPL doesn't have explicit syntax to call LITERAL_AGG. I will add comment to elaborate.

}
Integer dedupNumber = literal.getValueAs(Integer.class);
TopHitsAggregationBuilder topHitsAggregationBuilder =
AggregationBuilders.topHits(aggFieldName).from(0).fetchSource(false).size(dedupNumber);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you elaborate more on why we set fetchSource(false) here? The default value is true?

@LantaoJinLantaoJinNov 24, 2025

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Update: change to use fetchSource instead of fetchField due to fetchField cannot work on OS object type.

// LinkedHashMap["name" -> "A", "category" -> "X"],
// LinkedHashMap["name" -> "A", "category" -> "X"]
// ]
List<Map<String, Object>> res =

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket? I'm thinking Is it still appropriate to translate dedup x to aggregate? Normally speaking, aggregate should return only 1 row value for the metric in each bucket. While dedup x can return more than 1 rows and it now affects the API of agg metric parser.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

[question] Is this the only case that agg metric parser should return a List of value for each bucket?

Yeah, just top_hits metric agg does. top_hits metric agg in OpenSearch can be used in any bucket agg in DSL as a sub-aggregation, it's quite different with SQL' aggregate.

Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Signed-off-by: Lantao Jin <ltjin@amazon.com>
@@ -1601,15 +1600,15 @@ private static void buildDedupNotNull(
.partitionBy(dedupeFields)
.orderBy(dedupeFields)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Will dedupe work without ordering by deduped fields?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

em, it should work in non-pushdown case. maybe we could remove this orderBy in window.

@yuancuyuancuNov 26, 2025

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

It looks a little strange to me because the sort keys in a window should be the same (the partition key)

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Yes. it was not introduced by this pr. let's fix them in followup PR since it will change the plan a lot.


@ToString.Exclude private final Settings settings;

@ToString.Exclude private boolean topHitsAgg = false;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

What is this field used for? It seems it's not referred anywhere

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

yes. can be deleted in current impl

Comment on lines +1820 to +1828
@Ignore("https://github.com/opensearch-project/sql/issues/4789")
public void testDedupExpr() throws IOException {
enabledOnlyWhenPushdownIsEnabled();
String expected = loadExpectedPlan("explain_dedup_expr1.yaml");
assertYamlEqualsIgnoreId(
expected,
explainQueryYaml(
"source=opensearch-sql_test_index_account | eval new_gender = lower(gender) | dedup 1"
+ " new_gender"));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Isn't such case covered by CalcitePPLDedupIT.testDedupExpr

String.format("Unsupported push-down aggregator %s", aggCall.getAggregation()));
};
}
case LITERAL_AGG -> {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

One alternative in my mind is using a self-defined class extends calcite.Aggregate. It should also be able to leverage our AggregateAnalyzer.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

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

Can you check if we will never produce a LITERAL_AGG in RelBuilder or by some rules in planner?

Otherwise, we may push down a real LITERAL_AGG to be tophits while it shouldn't be.

Checked, so far no LITERAL_AGG can be produced because PPL doesn't support literal in aggregators.
In SQL

SELECT dept_id,
COUNT(*) as emp_count,
1 as constant_val
FROM employees
GROUP BY dept_id;

Could be rewritten to

SELECT dept_id,
COUNT(*) as emp_count,
LITERAL_AGG(1) as constant_val
FROM employees
GROUP BY dept_id;

to reduce a Project upon Aggregate.

From

Project(dept_id, emp_count, constant_val=1)
Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count)
Scan(employees)

To

Aggregate(GROUP BY dept_id, compute COUNT(*) as emp_count, compute LITERAL_AGG(1) as constant_val)
Scan(employees)

But stats only support pre-defined aggregators. cc @qianheng-aws

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

I see SubQueryRemoveRule will produce LITERAL_AGG for some(we don't have) and in subquery. But I'm not sure whether it will be triggered.

And RelBuilder's literalAgg method is public, we should avoid call that method by developers then.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

@LantaoJin
LantaoJin merged commit 5ceacb6 into opensearch-project:mainNov 26, 2025
39 checks passed
@opensearch-trigger-bot

Copy link
Copy Markdown
Contributor

The backport to 2.19-dev failed:

The process '/usr/bin/git' failed with exit code 128

To backport manually, run these commands in your terminal:

# Navigate to the root of your repositorycd$(git rev-parse --show-toplevel)# Fetch latest updates from GitHub
git fetch
# Create a new working tree
git worktree add ../.worktrees/sql/backport-2.19-dev 2.19-dev
# Navigate to the new working treepushd ../.worktrees/sql/backport-2.19-dev
# Create a new branch
git switch --create backport/backport-4844-to-2.19-dev
# Cherry-pick the merged commit of this pull request and resolve the conflicts
git cherry-pick -x --mainline 1 5ceacb6945d6bacc36e037ce4da215e1cf031b56
# Push it to GitHub
git push --set-upstream origin backport/backport-4844-to-2.19-dev
# Go back to the original working treepopd# Delete the working tree
git worktree remove ../.worktrees/sql/backport-2.19-dev

Then, create a pull request where the base branch is 2.19-dev and the compare/head branch is backport/backport-4844-to-2.19-dev.

LantaoJin added a commit to LantaoJin/search-plugins-sql that referenced this pull request Nov 27, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
(cherry picked from commit 5ceacb6)
@LantaoJinLantaoJin added the backport-manually Filed a PR to backport manually. label Nov 27, 2025
asifabashar pushed a commit to asifabashar/sql that referenced this pull request Dec 10, 2025
…4844)
* Enable dedup pushdown
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix doctest
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* refactor
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Disable dedup expr
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix IT
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix yaml test
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* add more comments in code
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* fix conflicts
Signed-off-by: Lantao Jin <ltjin@amazon.com>
* Address comments
Signed-off-by: Lantao Jin <ltjin@amazon.com>
---------
Signed-off-by: Lantao Jin <ltjin@amazon.com>
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

backport 2.19-devbackport-failedbackport-manuallyFiled a PR to backport manually.enhancementNew feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[ENHANCEMENT] Convert dedup pushdown to composite + top_hits

3 participants

@LantaoJin@yuancu@qianheng-aws