-
Notifications
You must be signed in to change notification settings - Fork 230
Pushdown case function in aggregations as range queries #4400
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from 25 commits
0acdd51
6afdcb6
f416dee
cbcb25a
7ec0684
51813f0
373e825
f3830a5
da5a9c5
b46dd95
0660994
0493629
f690add
b204b2e
4dc86db
d38a916
6beed21
d40c244
606e346
a5fdd66
fdb9886
e4c5266
7a8db58
0ca81aa
e701e57
4968d1c
a22ae79
010fd06
7d82cd8
3b65b2d
91aaee8
b8bb898
c950e16
0a5a55b
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -227,6 +227,16 @@ Argument type: all the supported data type, (NOTE : there is no comma before "el | |
|
|
||
| Return type: any | ||
|
|
||
| Limitations | ||
| >>>>>>>>>>> | ||
|
|
||
| When each condition is a field comparison with a numeric literal and each result expression is a string literal, the query will be optimized as `range aggregations <https://docs.opensearch.org/latest/aggregations/bucket/range>`_ if pushdown optimization is enabled. However, this optimization has the following limitations: | ||
|
|
||
| - Null values will not be grouped into any bucket of a range aggregation and will be ignored | ||
| - The default ELSE clause will use the string literal ``"null"`` instead of actual NULL values | ||
|
|
||
| To avoid these edge-case limitations, set ``plugins.calcite.pushdown.enabled`` to false. | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. remove this
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Removed |
||
|
|
||
| Example:: | ||
|
|
||
| os> source=accounts | eval result = case(age > 35, firstname, age < 30, lastname else employer) | fields result, firstname, lastname, age, employer | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -980,4 +980,111 @@ public void testFillNullValueSyntaxExplain() throws IOException { | |
| String.format( | ||
| "source=%s | fields age, balance | fillnull value=0", TEST_INDEX_ACCOUNT))); | ||
| } | ||
|
|
||
| @Test | ||
| public void testCasePushdownAsRangeQueryExplain() throws IOException { | ||
| // CASE 1: Range - Metric | ||
| // 1.1 Range - Metric | ||
| assertYamlEqualsJsonIgnoreId( | ||
| loadExpectedPlan("agg_range_metric_push.yaml"), | ||
| explainQueryToString( | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. should change to
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Fixed. Thanks for reminding! |
||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 'u30', age < 40, 'u40' else 'u100') |" | ||
| + " stats avg(age) as avg_age by age_range", | ||
| TEST_INDEX_BANK))); | ||
|
|
||
| // 1.2 Range - Metric (COUNT) | ||
| assertYamlEqualsJsonIgnoreId( | ||
| loadExpectedPlan("agg_range_count_push.yaml"), | ||
| explainQueryToString( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 'u30', age >= 30 and age < 40, 'u40'" | ||
| + " else 'u100') | stats avg(age) by age_range", | ||
| TEST_INDEX_BANK))); | ||
|
|
||
| // 1.3 Range - Range - Metric | ||
| assertYamlEqualsJsonIgnoreId( | ||
| loadExpectedPlan("agg_range_range_metric_push.yaml"), | ||
| explainQueryToString( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 'u30', age < 40, 'u40' else 'u100')," | ||
| + " balance_range = case(balance < 20000, 'medium' else 'high') | stats" | ||
| + " avg(balance) as avg_balance by age_range, balance_range", | ||
| TEST_INDEX_BANK))); | ||
|
|
||
| // 1.4 Range - Metric (With null & discontinuous ranges) | ||
| assertYamlEqualsJsonIgnoreId( | ||
| loadExpectedPlan("agg_range_metric_complex_push.yaml"), | ||
| explainQueryToString( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 'u30', (age >= 35 and age < 40) or age" | ||
| + " >= 80, '30-40 or >=80') | stats avg(balance) by age_range", | ||
| TEST_INDEX_BANK))); | ||
|
|
||
| // 1.5 Should not be pushed because the range is not closed-open | ||
| assertYamlEqualsJsonIgnoreId( | ||
| loadExpectedPlan("agg_case_cannot_push.yaml"), | ||
| explainQueryToString( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 'u30', age >= 30 and age <= 40, 'u40'" | ||
| + " else 'u100') | stats avg(age) as avg_age by age_range", | ||
| TEST_INDEX_BANK))); | ||
|
|
||
| // 1.6 Should not be pushed as range query because the result expression is not a string | ||
| // literal. | ||
| // Range aggregation keys must be strings | ||
| assertYamlEqualsJsonIgnoreId( | ||
| loadExpectedPlan("agg_case_num_res_cannot_push.yaml"), | ||
| explainQueryToString( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 30 else 100) | stats count() by" | ||
| + " age_range", | ||
| TEST_INDEX_BANK))); | ||
|
|
||
| // CASE 2: Composite - Range - Metric | ||
| // 2.1 Composite (term) - Range - Metric | ||
| assertYamlEqualsJsonIgnoreId( | ||
| loadExpectedPlan("agg_composite_range_metric_push.yaml"), | ||
| explainQueryToString( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 'u30' else 'a30') | stats avg(balance)" | ||
| + " by state, age_range", | ||
| TEST_INDEX_BANK))); | ||
|
|
||
| // 2.2 Composite (date histogram) - Range - Metric | ||
| assertYamlEqualsJsonIgnoreId( | ||
| loadExpectedPlan("agg_composite_date_range_push.yaml"), | ||
| explainQueryToString( | ||
| "source=opensearch-sql_test_index_time_data | eval value_range = case(value < 7000," | ||
| + " 'small' else 'large') | stats avg(value) by value_range, span(@timestamp," | ||
| + " 1h)")); | ||
|
|
||
| // 2.3 Composite(2 fields) - Range - Metric (with count) | ||
| assertYamlEqualsJsonIgnoreId( | ||
| loadExpectedPlan("agg_composite2_range_count_push.yaml"), | ||
| explainQueryToString( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 'u30' else 'a30') | stats" | ||
| + " avg(balance), count() by age_range, state, gender", | ||
| TEST_INDEX_BANK))); | ||
|
|
||
| // 2.4 Composite (2 fields) - Range - Range - Metric (with count) | ||
| assertYamlEqualsJsonIgnoreId( | ||
| loadExpectedPlan("agg_composite2_range_range_count_push.yaml"), | ||
| explainQueryToString( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 35, 'u35' else 'a35'), balance_range =" | ||
| + " case(balance < 20000, 'medium' else 'high') | stats avg(balance) as" | ||
| + " avg_balance by age_range, balance_range, state", | ||
| TEST_INDEX_BANK))); | ||
|
|
||
| // 2.5 Should not be pushed down as range query because case result expression is not constant | ||
| assertYamlEqualsJsonIgnoreId( | ||
| loadExpectedPlan("agg_case_composite_cannot_push.yaml"), | ||
| explainQueryToString( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 35, 'u35' else email) | stats avg(balance)" | ||
| + " as avg_balance by age_range, state", | ||
| TEST_INDEX_BANK))); | ||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -5,7 +5,10 @@ | |
|
|
||
| package org.opensearch.sql.calcite.remote; | ||
|
|
||
| import static org.opensearch.sql.legacy.TestsConstants.TEST_INDEX_BANK; | ||
| import static org.opensearch.sql.legacy.TestsConstants.TEST_INDEX_STATE_COUNTRY_WITH_NULL; | ||
| import static org.opensearch.sql.legacy.TestsConstants.TEST_INDEX_WEBLOGS; | ||
| import static org.opensearch.sql.util.MatcherUtils.closeTo; | ||
| import static org.opensearch.sql.util.MatcherUtils.rows; | ||
| import static org.opensearch.sql.util.MatcherUtils.schema; | ||
| import static org.opensearch.sql.util.MatcherUtils.verifyDataRows; | ||
|
|
@@ -25,6 +28,9 @@ public void init() throws Exception { | |
| enableCalcite(); | ||
|
|
||
| loadIndex(Index.WEBLOG); | ||
| loadIndex(Index.TIME_TEST_DATA); | ||
| loadIndex(Index.STATE_COUNTRY_WITH_NULL); | ||
| loadIndex(Index.BANK); | ||
| appendDataForBadResponse(); | ||
| } | ||
|
|
||
|
|
@@ -246,4 +252,222 @@ public void testCaseWhenInSubquery() throws IOException { | |
| rows("0.0.0.2", "GET", null, "4085", "500", "/shuttle/missions/sts-73/mission-sts-73.html"), | ||
| rows("::3", "GET", null, "3985", "403", "/shuttle/countdown/countdown.html")); | ||
| } | ||
|
|
||
| @Test | ||
| public void testCaseCanBePushedDownAsRangeQuery() throws IOException { | ||
| // CASE 1: Range - Metric | ||
| // 1.1 Range - Metric | ||
| JSONObject actual1 = | ||
| executeQuery( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 'u30', age < 40, 'u40' else 'u100') |" | ||
| + " stats avg(age) as avg_age by age_range", | ||
| TEST_INDEX_BANK)); | ||
| verifySchema(actual1, schema("avg_age", "double"), schema("age_range", "string")); | ||
| verifyDataRows(actual1, rows(28.0, "u30"), rows(35.0, "u40")); | ||
|
|
||
| // 1.2 Range - Metric (COUNT) | ||
| JSONObject actual2 = | ||
| executeQuery( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 'u30', age >= 30 and age < 40, 'u40'" | ||
| + " else 'u100') | stats avg(age) by age_range", | ||
| TEST_INDEX_BANK)); | ||
| verifySchema(actual2, schema("avg(age)", "double"), schema("age_range", "string")); | ||
| verifyDataRows(actual2, rows(28.0, "u30"), rows(35.0, "u40")); | ||
|
|
||
| // 1.3 Range - Range - Metric | ||
| JSONObject actual3 = | ||
| executeQuery( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 'u30', age < 40, 'u40' else 'u100')," | ||
| + " balance_range = case(balance < 20000, 'medium' else 'high') | stats" | ||
| + " avg(balance) as avg_balance by age_range, balance_range", | ||
| TEST_INDEX_BANK)); | ||
| verifySchema( | ||
| actual3, | ||
| schema("avg_balance", "double"), | ||
| schema("age_range", "string"), | ||
| schema("balance_range", "string")); | ||
| verifyDataRows( | ||
| actual3, | ||
| rows(32838.0, "u30", "high"), | ||
| closeTo(8761.333333333334, "u40", "medium"), | ||
| rows(42617.0, "u40", "high")); | ||
|
|
||
| // 1.4 Range - Metric (With null & discontinuous ranges) | ||
| JSONObject actual4 = | ||
| executeQuery( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 'u30', (age >= 35 and age < 40) or age" | ||
| + " >= 80, '30-40 or >=80') | stats avg(balance) by age_range", | ||
| TEST_INDEX_BANK)); | ||
| verifySchema(actual4, schema("avg(balance)", "double"), schema("age_range", "string")); | ||
| // There's such a discrepancy because null cannot be the key for a range query | ||
| if (isPushdownDisabled()) { | ||
| verifyDataRows( | ||
| actual4, | ||
| rows(32838.0, "u30"), | ||
| rows(30497.0, null), | ||
| closeTo(20881.333333333332, "30-40 or >=80")); | ||
| } else { | ||
| verifyDataRows( | ||
| actual4, | ||
| rows(32838.0, "u30"), | ||
| rows(30497.0, "null"), | ||
| closeTo(20881.333333333332, "30-40 or >=80")); | ||
| } | ||
|
|
||
| // 1.5 Should not be pushed because the range is not closed-open | ||
| JSONObject actual5 = | ||
| executeQuery( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 'u30', age >= 30 and age <= 40, 'u40'" | ||
| + " else 'u100') | stats avg(age) as avg_age by age_range", | ||
| TEST_INDEX_BANK)); | ||
| verifySchema(actual5, schema("avg_age", "double"), schema("age_range", "string")); | ||
| verifyDataRows(actual5, rows(35.0, "u40"), rows(28.0, "u30")); | ||
| } | ||
|
|
||
| @Test | ||
| public void testCaseCanBePushedDownAsCompositeRangeQuery() throws IOException { | ||
| // CASE 2: Composite - Range - Metric | ||
| // 2.1 Composite (term) - Range - Metric | ||
| JSONObject actual6 = | ||
| executeQuery( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 'u30' else 'a30') | stats avg(balance)" | ||
| + " by state, age_range", | ||
| TEST_INDEX_BANK)); | ||
| verifySchema( | ||
| actual6, | ||
| schema("avg(balance)", "double"), | ||
| schema("state", "string"), | ||
| schema("age_range", "string")); | ||
| verifyDataRows( | ||
| actual6, | ||
| rows(39225.0, "IL", "a30"), | ||
| rows(48086.0, "IN", "a30"), | ||
| rows(4180.0, "MD", "a30"), | ||
| rows(40540.0, "PA", "a30"), | ||
| rows(5686.0, "TN", "a30"), | ||
| rows(32838.0, "VA", "u30"), | ||
| rows(16418.0, "WA", "a30")); | ||
|
|
||
| // 2.2 Composite (date histogram) - Range - Metric | ||
| JSONObject actual7 = | ||
| executeQuery( | ||
| "source=opensearch-sql_test_index_time_data | eval value_range = case(value < 7000," | ||
| + " 'small' else 'large') | stats avg(value) by value_range, span(@timestamp," | ||
| + " 1month)"); | ||
| verifySchema( | ||
| actual7, | ||
| schema("avg(value)", "double"), | ||
| schema("span(@timestamp,1month)", "timestamp"), | ||
| schema("value_range", "string")); | ||
|
|
||
| verifyDataRows( | ||
| actual7, | ||
| closeTo(6642.521739130435, "2025-07-01 00:00:00", "small"), | ||
| closeTo(8381.917808219177, "2025-07-01 00:00:00", "large"), | ||
| rows(6489.0, "2025-08-01 00:00:00", "small"), | ||
| rows(8375.0, "2025-08-01 00:00:00", "large")); | ||
|
|
||
| // 2.3 Composite(2 fields) - Range - Metric (with count) | ||
| JSONObject actual8 = | ||
| executeQuery( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 30, 'u30' else 'a30') | stats" | ||
| + " avg(balance), count() by age_range, state, gender", | ||
| TEST_INDEX_BANK)); | ||
| verifySchema( | ||
| actual8, | ||
| schema("avg(balance)", "double"), | ||
| schema("count()", "bigint"), | ||
| schema("age_range", "string"), | ||
| schema("state", "string"), | ||
| schema("gender", "string")); | ||
| verifyDataRows( | ||
| actual8, | ||
| rows(5686.0, 1, "a30", "TN", "M"), | ||
| rows(16418.0, 1, "a30", "WA", "M"), | ||
| rows(40540.0, 1, "a30", "PA", "F"), | ||
| rows(4180.0, 1, "a30", "MD", "M"), | ||
| rows(32838.0, 1, "u30", "VA", "F"), | ||
| rows(39225.0, 1, "a30", "IL", "M"), | ||
| rows(48086.0, 1, "a30", "IN", "F")); | ||
|
|
||
| // 2.4 Composite (2 fields) - Range - Range - Metric (with count) | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Please add test case for
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Test added |
||
| JSONObject actual9 = | ||
| executeQuery( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 35, 'u35' else 'a35'), balance_range =" | ||
| + " case(balance < 20000, 'medium' else 'high') | stats avg(balance) as" | ||
| + " avg_balance by age_range, balance_range, state", | ||
| TEST_INDEX_BANK)); | ||
| verifySchema( | ||
| actual9, | ||
| schema("avg_balance", "double"), | ||
| schema("age_range", "string"), | ||
| schema("balance_range", "string"), | ||
| schema("state", "string")); | ||
| verifyDataRows( | ||
| actual9, | ||
| rows(39225.0, "u35", "high", "IL"), | ||
| rows(48086.0, "u35", "high", "IN"), | ||
| rows(4180.0, "u35", "medium", "MD"), | ||
| rows(40540.0, "a35", "high", "PA"), | ||
| rows(5686.0, "a35", "medium", "TN"), | ||
| rows(32838.0, "u35", "high", "VA"), | ||
| rows(16418.0, "a35", "medium", "WA")); | ||
|
|
||
| // 2.5 Should not be pushed because case result expression is not constant | ||
| JSONObject actual10 = | ||
| executeQuery( | ||
| String.format( | ||
| "source=%s | eval age_range = case(age < 35, 'u35' else email) | stats avg(balance)" | ||
| + " as avg_balance by age_range, state", | ||
| TEST_INDEX_BANK)); | ||
| verifySchema( | ||
| actual10, | ||
| schema("avg_balance", "double"), | ||
| schema("age_range", "string"), | ||
| schema("state", "string")); | ||
| verifyDataRows( | ||
| actual10, | ||
| rows(32838.0, "u35", "VA"), | ||
| rows(4180.0, "u35", "MD"), | ||
| rows(48086.0, "u35", "IN"), | ||
| rows(40540.0, "virginiaayala@filodyne.com", "PA"), | ||
| rows(39225.0, "u35", "IL"), | ||
| rows(5686.0, "hattiebond@netagy.com", "TN"), | ||
| rows(16418.0, "elinorratliff@scentric.com", "WA")); | ||
| } | ||
|
|
||
| @Test | ||
| public void testCaseAggWithNullValues() throws IOException { | ||
| JSONObject actual = | ||
| executeQuery( | ||
| String.format( | ||
| "source=%s" | ||
| + "| eval age_category = case(" | ||
| + " age < 20, 'teenager'," | ||
| + " age < 70, 'adult'," | ||
| + " age >= 70, 'senior'" | ||
| + " else 'unknown')" | ||
| + "| stats avg(age) by age_category", | ||
| TEST_INDEX_STATE_COUNTRY_WITH_NULL)); | ||
| verifySchema(actual, schema("avg(age)", "double"), schema("age_category", "string")); | ||
| // There is such discrepancy because range aggregations will ignore null values | ||
| if (isPushdownDisabled()) { | ||
| verifyDataRows( | ||
| actual, | ||
| rows(10, "teenager"), | ||
| rows(25, "adult"), | ||
| rows(70, "senior"), | ||
| rows(null, "unknown")); | ||
| } else { | ||
| verifyDataRows(actual, rows(10, "teenager"), rows(25, "adult"), rows(70, "senior")); | ||
| } | ||
| } | ||
| } | ||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
IMO, it's not a limitation of
casefunction, it is just a restricted optimization. We can just call out in what case, thecasefunction would be optimized to range DSL. Can we add some optimizablecaseusages in user doc?There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I think the stated conditions are in the scope of restricted optimizations, but the limitations are not because we will still do the optimization regardless of whether it has null values in its column or whether there is a default NULL range.
The problem is that there is no way to know in advance whether there exists null values in a column. Therefore, if we do this optimization, we always risk the discrepancy in results of with & without push-down.