This is an automated email from the ASF dual-hosted git repository.
yuqi1129 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/main by this push:
new b7add29edb [#11672] improvement(optimizer): optimize trigger-expr
evaluation with table metadata & partition short-circuit (#11664)
b7add29edb is described below
commit b7add29edbc5506e3ccff98884bb7081b327b009
Author: roryqi <[email protected]>
AuthorDate: Mon Jun 29 05:24:44 2026 -0400
[#11672] improvement(optimizer): optimize trigger-expr evaluation with
table metadata & partition short-circuit (#11664)
### What changes were proposed in this pull request?
Two related optimizations to the `maintenance/optimizer` recommender's
trigger-expr / score-expr evaluation:
1. **Extend evaluation context with table metadata** — the trigger-expr
context now exposes `column_count`, `partition_count`,
`sort_order_count`, and table properties (numeric values parsed to
`long`, others kept as `string`), in addition to partition and table
statistics. Both partitioned and non-partitioned tables now evaluate
against partition statistics (when present), table statistics, and table
metadata. The trigger-expr string representation is unchanged.
2. **Speed up partitioned-table evaluation** (port of Pinterest
gravitino-pinterest#249):
- Short-circuit: evaluate the expression with table-level context only;
if it resolves without referencing partition variables, skip the
per-partition loop (relies on the QL engine's left-to-right `&&` / `||`
short-circuiting).
- Precompute the table-level context once per `initialize()` instead of
rebuilding it for every partition.
- Cache compiled hyphen-to-underscore regex patterns in
`QLExpressionEvaluator` to avoid `Pattern.compile` on every evaluation.
- Adds `ExpressionEvaluator#tryToEvaluateBool` returning
`Optional<Boolean>`.
### Why are the changes needed?
Trigger expressions previously could only reference partition/table
statistics, limiting the rules users can write. They also re-evaluated
every partition even when a table-level expression already decided the
outcome, which is costly for large partitioned tables.
Fixes #11672
### Does this PR introduce _any_ user-facing change?
No API changes. Trigger-expr authors gain new referenceable variables
(`column_count`, `partition_count`, `sort_order_count`, and table
properties).
### How was this patch tested?
New/extended unit tests: `TestTableMetadataTriggerExpressionUtils`,
`TestQLExpressionEvaluator`, and `TestCompactionStrategyHandler`.
`./gradlew :maintenance:optimizer:test` passes locally.
---
.../handler/BaseExpressionStrategyHandler.java | 74 +++++-
.../recommender/util/ExpressionEvaluator.java | 15 ++
.../recommender/util/QLExpressionEvaluator.java | 81 ++++++-
.../util/TableMetadataTriggerExpressionUtils.java | 90 +++++++
.../compaction/TestCompactionStrategyHandler.java | 264 +++++++++++++++++++++
.../util/TestQLExpressionEvaluator.java | 65 +++++
.../TestTableMetadataTriggerExpressionUtils.java | 118 +++++++++
7 files changed, 686 insertions(+), 21 deletions(-)
diff --git
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/handler/BaseExpressionStrategyHandler.java
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/handler/BaseExpressionStrategyHandler.java
index 5c0aa59663..07072adfda 100644
---
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/handler/BaseExpressionStrategyHandler.java
+++
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/handler/BaseExpressionStrategyHandler.java
@@ -24,6 +24,7 @@ import java.util.Comparator;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
+import java.util.Optional;
import java.util.PriorityQueue;
import java.util.stream.Collectors;
import lombok.Value;
@@ -42,6 +43,7 @@ import
org.apache.gravitino.maintenance.optimizer.recommender.util.ExpressionEva
import
org.apache.gravitino.maintenance.optimizer.recommender.util.QLExpressionEvaluator;
import
org.apache.gravitino.maintenance.optimizer.recommender.util.StatisticsUtils;
import
org.apache.gravitino.maintenance.optimizer.recommender.util.StrategyUtils;
+import
org.apache.gravitino.maintenance.optimizer.recommender.util.TableMetadataTriggerExpressionUtils;
import org.apache.gravitino.rel.Table;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -65,6 +67,8 @@ public abstract class BaseExpressionStrategyHandler
implements StrategyHandler {
private Map<PartitionPath, List<StatisticEntry<?>>> partitionStatistics;
private Table tableMetadata;
private NameIdentifier nameIdentifier;
+ // Cached table-level context (table stats + metadata + rules), computed
once per initialize().
+ private Map<String, Object> tableLevelContext;
/** Create a handler that evaluates expressions with the default QL
evaluator. */
protected BaseExpressionStrategyHandler() {
@@ -79,6 +83,7 @@ public abstract class BaseExpressionStrategyHandler
implements StrategyHandler {
this.strategy = context.strategy();
this.tableStatistics = context.tableStatistics();
this.partitionStatistics = context.partitionStatistics();
+ this.tableLevelContext = buildTableLevelContext();
}
@Override
@@ -131,16 +136,30 @@ public abstract class BaseExpressionStrategyHandler
implements StrategyHandler {
return false;
}
String triggerExpression = triggerExpression(strategy);
- return partitionStatistics.values().stream()
- .anyMatch(partitionStats -> evaluateBool(triggerExpression,
partitionStats));
+
+ // Short-circuit: try to evaluate the trigger expression with table-level
context only
+ // (no partition statistics). This relies on the QL engine's left-to-right
short-circuit
+ // evaluation of && operators: if a table-level predicate appears first
and evaluates to
+ // false, the engine returns false without resolving subsequent
partition-level variables.
+ // Note: this optimization only works when the table-level predicate is
positioned before the
+ // partition-level predicate, e.g. "sort_order_count > 0 && datafile_mse <
limit" will
+ // short-circuit, but "datafile_mse < limit && sort_order_count > 0" will
not.
+ Optional<Boolean> evaluationResultWithoutPartitions =
+ tryToEvaluateBool(triggerExpression, List.of());
+ // If the trigger expression can be evaluated with table-level variables
only, return the
+ // result. Otherwise, evaluate the trigger expression for each partition.
+ return evaluationResultWithoutPartitions.orElseGet(
+ () ->
+ partitionStatistics.values().stream()
+ .anyMatch(partitionStats -> evaluateBool(triggerExpression,
partitionStats)));
}
private boolean shouldTriggerForNonPartitionTable() {
- return evaluateBool(triggerExpression(strategy), tableStatistics);
+ return evaluateBool(triggerExpression(strategy), List.of());
}
private StrategyEvaluation evaluateForNonPartitionTable() {
- long score = evaluateLong(scoreExpression(strategy), tableStatistics);
+ long score = evaluateLong(scoreExpression(strategy), List.of());
if (score <= 0) {
return StrategyEvaluation.NO_EXECUTION;
}
@@ -200,7 +219,7 @@ public abstract class BaseExpressionStrategyHandler
implements StrategyHandler {
}
private long evaluateLong(String expression, List<StatisticEntry<?>>
statistics) {
- Map<String, Object> context = buildExpressionContext(strategy, statistics);
+ Map<String, Object> context = buildExpressionContext(statistics);
try {
return expressionEvaluator.evaluateLong(expression, context);
} catch (RuntimeException e) {
@@ -210,7 +229,7 @@ public abstract class BaseExpressionStrategyHandler
implements StrategyHandler {
}
private boolean evaluateBool(String expression, List<StatisticEntry<?>>
statistics) {
- Map<String, Object> context = buildExpressionContext(strategy, statistics);
+ Map<String, Object> context = buildExpressionContext(statistics);
try {
return expressionEvaluator.evaluateBool(expression, context);
} catch (RuntimeException e) {
@@ -219,10 +238,47 @@ public abstract class BaseExpressionStrategyHandler
implements StrategyHandler {
}
}
- private static Map<String, Object> buildExpressionContext(
- Strategy strategy, List<StatisticEntry<?>> statistics) {
- Map<String, Object> context = new HashMap<>();
+ private Optional<Boolean> tryToEvaluateBool(
+ String expression, List<StatisticEntry<?>> statistics) {
+ Map<String, Object> context = buildExpressionContext(statistics);
+ try {
+ return expressionEvaluator.tryToEvaluateBool(expression, context);
+ } catch (RuntimeException e) {
+ LOG.warn("Failed to evaluate expression '{}' with context {}",
expression, context, e);
+ // Per ExpressionEvaluator#tryToEvaluateBool, a failed evaluation yields
an empty Optional so
+ // the caller falls back to per-partition evaluation instead of treating
it as "do not
+ // trigger".
+ return Optional.empty();
+ }
+ }
+
+ /**
+ * Build a combined evaluation context. The {@code statistics} argument
varies per partition for
+ * partitioned tables; it is layered on top of the precomputed {@link
#tableLevelContext} (table
+ * statistics, table metadata, and numeric rule values) so that {@code
trigger-expr} and {@code
+ * score-expr} can reference all of them.
+ *
+ * @param statistics statistics of the unit being evaluated (a single
partition, or empty for the
+ * table-level evaluation)
+ * @return combined context
+ */
+ private Map<String, Object> buildExpressionContext(List<StatisticEntry<?>>
statistics) {
+ Map<String, Object> context = new HashMap<>(tableLevelContext);
context.putAll(StatisticsUtils.buildStatisticsContext(statistics));
+ return context;
+ }
+
+ /** Precompute the table-level context shared across all partition
evaluations. */
+ private Map<String, Object> buildTableLevelContext() {
+ Map<String, Object> context = new HashMap<>();
+ context.putAll(StatisticsUtils.buildStatisticsContext(tableStatistics));
+
context.putAll(TableMetadataTriggerExpressionUtils.buildTableMetadataContext(tableMetadata));
+ context.putAll(getRulesExpressionContext(strategy));
+ return context;
+ }
+
+ private static Map<String, Object> getRulesExpressionContext(Strategy
strategy) {
+ Map<String, Object> context = new HashMap<>();
strategy
.rules()
.forEach(
diff --git
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/util/ExpressionEvaluator.java
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/util/ExpressionEvaluator.java
index 71d8a4d6e1..fb6462b8e3 100644
---
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/util/ExpressionEvaluator.java
+++
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/util/ExpressionEvaluator.java
@@ -20,6 +20,7 @@
package org.apache.gravitino.maintenance.optimizer.recommender.util;
import java.util.Map;
+import java.util.Optional;
/**
* Evaluates rule expressions against a provided context map.
@@ -46,4 +47,18 @@ public interface ExpressionEvaluator {
* @return evaluation result as a {@code long}
*/
long evaluateLong(String expression, Map<String, Object> context);
+
+ /**
+ * Evaluates an expression that returns a boolean value and returns a {@link
java.util.Optional}
+ * containing the result if the evaluation is successful, or an empty {@link
java.util.Optional}
+ * if it fails.
+ *
+ * <p>Evaluation may fail if there are syntax errors in a specified
expression or if a variable
+ * referenced in the expression is not present in the context.
+ *
+ * @param expression expression to evaluate
+ * @param context variable bindings for the expression
+ * @return evaluation result as an {@link java.util.Optional} containing a
{@code boolean}
+ */
+ Optional<Boolean> tryToEvaluateBool(String expression, Map<String, Object>
context);
}
diff --git
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/util/QLExpressionEvaluator.java
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/util/QLExpressionEvaluator.java
index cd7609ebca..1a547aeac6 100644
---
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/util/QLExpressionEvaluator.java
+++
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/util/QLExpressionEvaluator.java
@@ -22,17 +22,31 @@ package
org.apache.gravitino.maintenance.optimizer.recommender.util;
import com.alibaba.qlexpress4.Express4Runner;
import com.alibaba.qlexpress4.InitOptions;
import com.alibaba.qlexpress4.QLOptions;
+import com.alibaba.qlexpress4.exception.QLException;
import com.google.common.base.Preconditions;
import java.math.BigDecimal;
import java.util.Map;
+import java.util.Optional;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import java.util.stream.Collectors;
import org.apache.commons.lang3.StringUtils;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
public class QLExpressionEvaluator implements ExpressionEvaluator {
+
+ private static final Logger LOG =
LoggerFactory.getLogger(QLExpressionEvaluator.class);
+
private static final Express4Runner RUNNER = new
Express4Runner(InitOptions.DEFAULT_OPTIONS);
+ // Cache compiled regex patterns for hyphen-to-underscore replacement, keyed
by the set of
+ // hyphenated context keys. Avoids expensive Pattern.compile on every
evaluation call.
+ private final Map<Set<String>, HyphenReplacementRule> replacementRuleCache =
+ new ConcurrentHashMap<>();
+
@Override
public long evaluateLong(String expression, Map<String, Object> context) {
return toLong(evaluate(expression, context));
@@ -43,6 +57,11 @@ public class QLExpressionEvaluator implements
ExpressionEvaluator {
return (boolean) evaluate(expression, context);
}
+ @Override
+ public Optional<Boolean> tryToEvaluateBool(String expression, Map<String,
Object> context) {
+ return tryToEvaluate(expression, context).map(o -> (Boolean) o);
+ }
+
private Object evaluate(String expression, Map<String, Object> context) {
Preconditions.checkArgument(StringUtils.isNotBlank(expression),
"expression is blank");
Preconditions.checkArgument(context != null, "context is null");
@@ -52,6 +71,16 @@ public class QLExpressionEvaluator implements
ExpressionEvaluator {
.getResult();
}
+ private Optional<Object> tryToEvaluate(String expression, Map<String,
Object> context) {
+ try {
+ Object result = evaluate(expression, context);
+ return Optional.of(result);
+ } catch (QLException e) {
+ LOG.warn("Failed to evaluate expression '{}': {}", expression,
e.getMessage());
+ return Optional.empty();
+ }
+ }
+
private Map<String, Object> formatContextKey(Map<String, Object> context) {
return context.entrySet().stream()
.collect(
@@ -60,29 +89,57 @@ public class QLExpressionEvaluator implements
ExpressionEvaluator {
}
private String formatExpression(String expression, Map<String, Object>
context) {
- Map<String, String> replacements =
- context.keySet().stream()
- .collect(
- Collectors.toMap(key -> key, this::normalizeIdentifier, (left,
right) -> left));
- replacements.entrySet().removeIf(entry ->
entry.getKey().equals(entry.getValue()));
- if (replacements.isEmpty()) {
+ HyphenReplacementRule rule = getReplacementRule(context);
+ if (rule == null) {
return expression;
}
-
- String alternation =
-
replacements.keySet().stream().map(Pattern::quote).collect(Collectors.joining("|"));
- Pattern pattern = Pattern.compile("(?<![A-Za-z0-9_])(" + alternation +
")(?![A-Za-z0-9_])");
- Matcher matcher = pattern.matcher(expression);
+ Matcher matcher = rule.pattern.matcher(expression);
StringBuffer buffer = new StringBuffer();
while (matcher.find()) {
String matched = matcher.group(1);
matcher.appendReplacement(
- buffer, Matcher.quoteReplacement(replacements.getOrDefault(matched,
matched)));
+ buffer,
Matcher.quoteReplacement(rule.replacements.getOrDefault(matched, matched)));
}
matcher.appendTail(buffer);
return buffer.toString();
}
+ /** Return a cached replacement rule for the hyphenated keys in the context,
or null if none. */
+ private HyphenReplacementRule getReplacementRule(Map<String, Object>
context) {
+ Set<String> hyphenatedKeys =
+ context.keySet().stream()
+ .filter(key -> !key.equals(normalizeIdentifier(key)))
+ .collect(Collectors.toSet());
+ if (hyphenatedKeys.isEmpty()) {
+ return null;
+ }
+ return replacementRuleCache.computeIfAbsent(
+ hyphenatedKeys,
+ keys -> {
+ Map<String, String> replacements =
+ keys.stream()
+ .collect(
+ Collectors.toMap(
+ key -> key, this::normalizeIdentifier, (left, right)
-> left));
+ String alternation =
+
replacements.keySet().stream().map(Pattern::quote).collect(Collectors.joining("|"));
+ Pattern pattern =
+ Pattern.compile("(?<![A-Za-z0-9_])(" + alternation +
")(?![A-Za-z0-9_])");
+ return new HyphenReplacementRule(pattern, replacements);
+ });
+ }
+
+ /** Pre-compiled pattern and replacement map for rewriting hyphenated
identifiers. */
+ private static class HyphenReplacementRule {
+ final Pattern pattern;
+ final Map<String, String> replacements;
+
+ HyphenReplacementRule(Pattern pattern, Map<String, String> replacements) {
+ this.pattern = pattern;
+ this.replacements = replacements;
+ }
+ }
+
private String normalizeIdentifier(String name) {
return name.replace("-", "_");
}
diff --git
a/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/util/TableMetadataTriggerExpressionUtils.java
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/util/TableMetadataTriggerExpressionUtils.java
new file mode 100644
index 0000000000..24e539518e
--- /dev/null
+++
b/maintenance/optimizer/src/main/java/org/apache/gravitino/maintenance/optimizer/recommender/util/TableMetadataTriggerExpressionUtils.java
@@ -0,0 +1,90 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.gravitino.maintenance.optimizer.recommender.util;
+
+import com.google.common.base.Preconditions;
+import java.util.HashMap;
+import java.util.Map;
+import org.apache.commons.lang3.math.NumberUtils;
+import org.apache.gravitino.rel.Table;
+
+/**
+ * Utility for translating table metadata into variables usable in a {@code
trigger-expr} or {@code
+ * score-expr} expression.
+ */
+public class TableMetadataTriggerExpressionUtils {
+
+ private TableMetadataTriggerExpressionUtils() {}
+
+ /**
+ * Builds a context map for evaluating table metadata trigger expressions.
+ *
+ * <p>This method creates a map containing:
+ *
+ * <ul>
+ * <li>Built-in table metadata: {@code column_count}, {@code
partition_count}, {@code
+ * sort_order_count}
+ * <li>All table properties from {@link Table#properties()}, with numeric
strings converted to
+ * {@code long} values and all others kept as {@code String} values
+ * </ul>
+ *
+ * @param tableMetadata the table metadata to extract properties from
+ * @return a context map suitable for expression evaluation
+ */
+ public static Map<String, Object> buildTableMetadataContext(Table
tableMetadata) {
+ Preconditions.checkArgument(tableMetadata != null, "Table metadata is
null");
+ Map<String, Object> context = new HashMap<>();
+ context.put(
+ TableMetadataTriggerExpression.COLUMN_COUNT.getName(),
+ tableMetadata.columns() != null ? tableMetadata.columns().length : 0);
+ context.put(
+ TableMetadataTriggerExpression.PARTITION_COUNT.getName(),
+ tableMetadata.partitioning() != null ?
tableMetadata.partitioning().length : 0);
+ context.put(
+ TableMetadataTriggerExpression.SORT_ORDER_COUNT.getName(),
+ tableMetadata.sortOrder() != null ? tableMetadata.sortOrder().length :
0);
+
+ if (tableMetadata.properties() != null) {
+ tableMetadata
+ .properties()
+ .forEach(
+ (k, v) -> {
+ if (NumberUtils.isCreatable(v)) {
+ context.put(k, NumberUtils.createNumber(v).longValue());
+ } else {
+ // For non-numeric properties, we just put the value as is.
+ context.put(k, v);
+ }
+ });
+ }
+ return context;
+ }
+
+ /** Supported built-in table metadata trigger expression variables. */
+ private enum TableMetadataTriggerExpression {
+ COLUMN_COUNT,
+ PARTITION_COUNT,
+ SORT_ORDER_COUNT;
+
+ public String getName() {
+ return name().toLowerCase();
+ }
+ }
+}
diff --git
a/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/recommender/handler/compaction/TestCompactionStrategyHandler.java
b/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/recommender/handler/compaction/TestCompactionStrategyHandler.java
index 231204b696..300949a679 100644
---
a/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/recommender/handler/compaction/TestCompactionStrategyHandler.java
+++
b/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/recommender/handler/compaction/TestCompactionStrategyHandler.java
@@ -39,6 +39,9 @@ import
org.apache.gravitino.maintenance.optimizer.recommender.strategy.Gravitino
import
org.apache.gravitino.maintenance.optimizer.recommender.util.StrategyUtils;
import org.apache.gravitino.rel.Column;
import org.apache.gravitino.rel.Table;
+import org.apache.gravitino.rel.expressions.NamedReference;
+import org.apache.gravitino.rel.expressions.sorts.SortOrder;
+import org.apache.gravitino.rel.expressions.sorts.SortOrders;
import org.apache.gravitino.rel.expressions.transforms.Transforms;
import org.apache.gravitino.stats.StatisticValues;
import org.junit.jupiter.api.Assertions;
@@ -124,6 +127,228 @@ class TestCompactionStrategyHandler {
Assertions.assertFalse(lowHandler.shouldTrigger());
}
+ @Test
+ void testShouldTriggerUsingTableMetadata() {
+ NameIdentifier tableId = NameIdentifier.of("db", "table");
+ Map<String, Object> rules = new HashMap<>();
+ rules.put(StrategyUtils.TRIGGER_EXPR, "sort_order_count > 0 &&
min_target_size > 100");
+ rules.put(StrategyUtils.SCORE_EXPR, "1");
+ Strategy metadataStrategy =
+ new Strategy() {
+ @Override
+ public String name() {
+ return "metadata-trigger-test";
+ }
+
+ @Override
+ public String strategyType() {
+ return CompactionStrategyHandler.NAME;
+ }
+
+ @Override
+ public Map<String, Object> rules() {
+ return rules;
+ }
+
+ @Override
+ public Map<String, String> properties() {
+ return Map.of();
+ }
+
+ @Override
+ public Map<String, String> jobOptions() {
+ return Map.of();
+ }
+
+ @Override
+ public String jobTemplateName() {
+ return "compaction-template";
+ }
+ };
+
+ // A sorted table whose property satisfies the threshold should trigger.
+ Table sorted = Mockito.mock(Table.class);
+ Mockito.when(sorted.partitioning())
+ .thenReturn(new
org.apache.gravitino.rel.expressions.transforms.Transform[0]);
+ Mockito.when(sorted.sortOrder())
+ .thenReturn(new
org.apache.gravitino.rel.expressions.sorts.SortOrder[1]);
+ Mockito.when(sorted.properties()).thenReturn(Map.of("min_target_size",
"500"));
+ StrategyHandlerContext sortedContext =
+ StrategyHandlerContext.builder(tableId, metadataStrategy)
+ .withTableMetadata(sorted)
+ .withTableStatistics(List.of())
+ .build();
+ CompactionStrategyHandler sortedHandler = new CompactionStrategyHandler();
+ sortedHandler.initialize(sortedContext);
+ Assertions.assertTrue(sortedHandler.shouldTrigger());
+
+ // An unsorted table (sort_order_count == 0) must not trigger the same
expression.
+ Table unsorted = Mockito.mock(Table.class);
+ Mockito.when(unsorted.partitioning())
+ .thenReturn(new
org.apache.gravitino.rel.expressions.transforms.Transform[0]);
+ Mockito.when(unsorted.sortOrder())
+ .thenReturn(new
org.apache.gravitino.rel.expressions.sorts.SortOrder[0]);
+ Mockito.when(unsorted.properties()).thenReturn(Map.of("min_target_size",
"500"));
+ StrategyHandlerContext unsortedContext =
+ StrategyHandlerContext.builder(tableId, metadataStrategy)
+ .withTableMetadata(unsorted)
+ .withTableStatistics(List.of())
+ .build();
+ CompactionStrategyHandler unsortedHandler = new
CompactionStrategyHandler();
+ unsortedHandler.initialize(unsortedContext);
+ Assertions.assertFalse(unsortedHandler.shouldTrigger());
+ }
+
+ @Test
+ void testShouldTriggerShortCircuitsWhenTableLevelExprEvaluatesToFalse() {
+ NameIdentifier tableId = NameIdentifier.of("db", "table");
+ Table tableMetadata = Mockito.mock(Table.class);
+ Mockito.when(tableMetadata.partitioning())
+ .thenReturn(
+ new org.apache.gravitino.rel.expressions.transforms.Transform[] {
+ Transforms.identity("table")
+ });
+ Mockito.when(tableMetadata.sortOrder()).thenReturn(new SortOrder[0]);
+ Mockito.when(tableMetadata.columns()).thenReturn(new Column[0]);
+ Mockito.when(tableMetadata.properties()).thenReturn(Map.of());
+
+ // Since sort_order_count == 0, QL evaluator short-circuits the && and
returns false
+ // without needing to resolve the partition-level variable
+ Strategy strategy = buildStrategyWithTriggerExpr("sort_order_count > 0 &&
datafile_mse > 0");
+
+ Map<PartitionPath, List<StatisticEntry<?>>> partitionStats =
+ Map.of(
+ PartitionPath.of(Arrays.asList(new PartitionEntryImpl("table",
"table_1"))),
+ List.of(new StatisticEntryImpl("datafile_mse",
StatisticValues.longValue(10L))),
+ PartitionPath.of(Arrays.asList(new PartitionEntryImpl("table",
"table_2"))),
+ List.of(new StatisticEntryImpl("datafile_mse",
StatisticValues.longValue(20L))));
+
+ StrategyHandlerContext context =
+ StrategyHandlerContext.builder(tableId, strategy)
+ .withTableMetadata(tableMetadata)
+ .withTableStatistics(List.of())
+ .withPartitionStatistics(partitionStats)
+ .build();
+
+ CompactionStrategyHandler handler = new CompactionStrategyHandler();
+ handler.initialize(context);
+ Assertions.assertFalse(
+ handler.shouldTrigger(),
+ "Should short-circuit to false when table-level predicate is first and
evaluates to false");
+ }
+
+ @Test
+ void testShouldTriggerShortCircuitsWhenTableLevelExprEvaluatesToTrue() {
+ NameIdentifier tableId = NameIdentifier.of("db", "table");
+ Table tableMetadata = Mockito.mock(Table.class);
+ Mockito.when(tableMetadata.partitioning())
+ .thenReturn(
+ new org.apache.gravitino.rel.expressions.transforms.Transform[] {
+ Transforms.identity("table")
+ });
+ Mockito.when(tableMetadata.sortOrder())
+ .thenReturn(new SortOrder[]
{SortOrders.ascending(NamedReference.field("db"))});
+ Mockito.when(tableMetadata.columns()).thenReturn(new Column[0]);
+ Mockito.when(tableMetadata.properties()).thenReturn(Map.of());
+
+ // Since sort_order_count == 1, QL evaluator short-circuits the || and
returns true
+ // without needing to resolve the partition-level variable
+ Strategy strategy = buildStrategyWithTriggerExpr("sort_order_count > 0 ||
datafile_mse == 0");
+
+ Map<PartitionPath, List<StatisticEntry<?>>> partitionStats =
+ Map.of(
+ PartitionPath.of(Arrays.asList(new PartitionEntryImpl("table",
"table_1"))),
+ List.of(new StatisticEntryImpl("datafile_mse",
StatisticValues.longValue(10L))),
+ PartitionPath.of(Arrays.asList(new PartitionEntryImpl("table",
"table_2"))),
+ List.of(new StatisticEntryImpl("datafile_mse",
StatisticValues.longValue(20L))));
+
+ StrategyHandlerContext context =
+ StrategyHandlerContext.builder(tableId, strategy)
+ .withTableMetadata(tableMetadata)
+ .withTableStatistics(List.of())
+ .withPartitionStatistics(partitionStats)
+ .build();
+
+ CompactionStrategyHandler handler = new CompactionStrategyHandler();
+ handler.initialize(context);
+ Assertions.assertTrue(
+ handler.shouldTrigger(),
+ "Should short-circuit to true when table-level predicate is first and
evaluates to true");
+ }
+
+ @Test
+ void testShouldTriggerShortCircuitsWhenTableLevelOnlyExpr() {
+ NameIdentifier tableId = NameIdentifier.of("db", "table");
+ Table tableMetadata = Mockito.mock(Table.class);
+ Mockito.when(tableMetadata.partitioning())
+ .thenReturn(
+ new org.apache.gravitino.rel.expressions.transforms.Transform[] {
+ Transforms.identity("p")
+ });
+ Mockito.when(tableMetadata.sortOrder())
+ .thenReturn(new SortOrder[]
{SortOrders.ascending(NamedReference.field("id"))});
+ Mockito.when(tableMetadata.columns()).thenReturn(new Column[0]);
+ Mockito.when(tableMetadata.properties()).thenReturn(Map.of());
+
+ // Expression uses only table-level variable, should evaluate once without
partition iteration
+ Strategy strategy = buildStrategyWithTriggerExpr("sort_order_count > 0");
+
+ Map<PartitionPath, List<StatisticEntry<?>>> partitionStats =
+ Map.of(
+ PartitionPath.of(Arrays.asList(new PartitionEntryImpl("p", "1"))),
+ List.of(new StatisticEntryImpl("datafile_mse",
StatisticValues.longValue(10L))));
+
+ StrategyHandlerContext context =
+ StrategyHandlerContext.builder(tableId, strategy)
+ .withTableMetadata(tableMetadata)
+ .withTableStatistics(List.of())
+ .withPartitionStatistics(partitionStats)
+ .build();
+
+ CompactionStrategyHandler handler = new CompactionStrategyHandler();
+ handler.initialize(context);
+ Assertions.assertTrue(
+ handler.shouldTrigger(),
+ "Should short-circuit to true when table-level-only expression
evaluates to true");
+ }
+
+ @Test
+ void testShouldTriggerFallsBackToPartitionEvalWhenTableLevelExprIsNotFirst()
{
+ NameIdentifier tableId = NameIdentifier.of("db", "table");
+ Table tableMetadata = Mockito.mock(Table.class);
+ Mockito.when(tableMetadata.partitioning())
+ .thenReturn(
+ new org.apache.gravitino.rel.expressions.transforms.Transform[] {
+ Transforms.identity("table")
+ });
+ Mockito.when(tableMetadata.sortOrder()).thenReturn(new SortOrder[0]);
+ Mockito.when(tableMetadata.columns()).thenReturn(new Column[0]);
+ Mockito.when(tableMetadata.properties()).thenReturn(Map.of());
+
+ // Partition-level predicate first: datafile_mse > 0 && sort_order_count > 0
+ // tryToEvaluateBool will fail on missing 'datafile_mse', so falls back to
per-partition eval
+ Strategy strategy = buildStrategyWithTriggerExpr("datafile_mse > 0 &&
sort_order_count > 0");
+
+ Map<PartitionPath, List<StatisticEntry<?>>> partitionStats =
+ Map.of(
+ PartitionPath.of(Arrays.asList(new PartitionEntryImpl("table",
"table_1"))),
+ List.of(new StatisticEntryImpl("datafile_mse",
StatisticValues.longValue(10L))));
+
+ StrategyHandlerContext context =
+ StrategyHandlerContext.builder(tableId, strategy)
+ .withTableMetadata(tableMetadata)
+ .withTableStatistics(List.of())
+ .withPartitionStatistics(partitionStats)
+ .build();
+
+ CompactionStrategyHandler handler = new CompactionStrategyHandler();
+ handler.initialize(context);
+ // sort_order_count == 0 so per-partition eval will also return false
+ Assertions.assertFalse(
+ handler.shouldTrigger(),
+ "Falls back to per-partition eval; still false because
sort_order_count is 0");
+ }
+
@Test
void testJobConfig() {
NameIdentifier tableId = NameIdentifier.of("db", "table");
@@ -417,6 +642,45 @@ class TestCompactionStrategyHandler {
return evaluation.score();
}
+ private Strategy buildStrategyWithTriggerExpr(String triggerExpr) {
+ Map<String, Object> rules = new HashMap<>();
+ if (triggerExpr != null) {
+ rules.put(StrategyUtils.TRIGGER_EXPR, triggerExpr);
+ }
+ rules.put(StrategyUtils.SCORE_EXPR, "datafile_mse");
+ return new Strategy() {
+ @Override
+ public String name() {
+ return "compaction-test";
+ }
+
+ @Override
+ public String strategyType() {
+ return CompactionStrategyHandler.NAME;
+ }
+
+ @Override
+ public Map<String, Object> rules() {
+ return rules;
+ }
+
+ @Override
+ public Map<String, String> properties() {
+ return Map.of();
+ }
+
+ @Override
+ public Map<String, String> jobOptions() {
+ return Map.of();
+ }
+
+ @Override
+ public String jobTemplateName() {
+ return "compaction-template";
+ }
+ };
+ }
+
private Strategy buildStrategy(ScoreMode scoreMode, Integer maxPartitionNum)
{
Map<String, Object> rules = new HashMap<>();
rules.put(StrategyUtils.TRIGGER_EXPR, "datafile_mse > 0");
diff --git
a/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/recommender/util/TestQLExpressionEvaluator.java
b/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/recommender/util/TestQLExpressionEvaluator.java
index cf372fcd38..7d7877ab98 100644
---
a/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/recommender/util/TestQLExpressionEvaluator.java
+++
b/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/recommender/util/TestQLExpressionEvaluator.java
@@ -25,6 +25,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
import java.util.HashMap;
import java.util.Map;
+import java.util.Optional;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
@@ -167,4 +168,68 @@ public class TestQLExpressionEvaluator {
long result = evaluator.evaluateLong("a + metric-1", context);
assertEquals(5L, result);
}
+
+ @Test
+ void testTryToEvaluateBoolWithResult() {
+ Map<String, Object> context = new HashMap<>();
+ context.put("x", 5);
+ context.put("y", 10);
+
+ Optional<Boolean> result = evaluator.tryToEvaluateBool("x < y", context);
+ assertEquals(Optional.of(true), result);
+ }
+
+ @Test
+ void testTryToEvaluateBoolWithMissingVariableReturnsEmpty() {
+ Map<String, Object> context = new HashMap<>();
+ context.put("x", 5);
+
+ Optional<Boolean> result = evaluator.tryToEvaluateBool("x < y", context);
+ assertEquals(Optional.empty(), result);
+ }
+
+ @Test
+ void testTryToEvaluateBoolWithBlankExpression() {
+ assertThrows(IllegalArgumentException.class, () ->
evaluator.tryToEvaluateBool("", Map.of()));
+ assertThrows(IllegalArgumentException.class, () ->
evaluator.tryToEvaluateBool(null, Map.of()));
+ }
+
+ @Test
+ void testRepeatedEvaluationsWithSameHyphenatedKeysProduceCorrectResults() {
+ Map<String, Object> context1 = new HashMap<>();
+ context1.put("metric-a", 10);
+ context1.put("threshold", 5);
+ assertTrue(evaluator.evaluateBool("metric-a > threshold", context1));
+
+ // Same hyphenated key set, different values — should reuse cached
replacement rule
+ Map<String, Object> context2 = new HashMap<>();
+ context2.put("metric-a", 3);
+ context2.put("threshold", 5);
+ Assertions.assertFalse(evaluator.evaluateBool("metric-a > threshold",
context2));
+ }
+
+ @Test
+ void testDifferentHyphenatedKeySetsProduceCorrectResults() {
+ Map<String, Object> contextA = new HashMap<>();
+ contextA.put("metric-a", 10);
+ contextA.put("limit", 5);
+ assertEquals(5L, evaluator.evaluateLong("metric-a - limit", contextA));
+
+ // Different hyphenated key set — must not reuse the previous cached rule
+ Map<String, Object> contextB = new HashMap<>();
+ contextB.put("metric-b", 20);
+ contextB.put("limit", 8);
+ assertEquals(12L, evaluator.evaluateLong("metric-b - limit", contextB));
+ }
+
+ @Test
+ void testContextWithoutHyphenatedKeysBypassesReplacementCache() {
+ Map<String, Object> context = new HashMap<>();
+ context.put("alpha", 7);
+ context.put("beta", 3);
+
+ // No hyphenated keys — cache should not be involved, expression evaluated
as-is
+ assertEquals(10L, evaluator.evaluateLong("alpha + beta", context));
+ assertTrue(evaluator.evaluateBool("alpha > beta", context));
+ }
}
diff --git
a/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/recommender/util/TestTableMetadataTriggerExpressionUtils.java
b/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/recommender/util/TestTableMetadataTriggerExpressionUtils.java
new file mode 100644
index 0000000000..0e0c94af2f
--- /dev/null
+++
b/maintenance/optimizer/src/test/java/org/apache/gravitino/maintenance/optimizer/recommender/util/TestTableMetadataTriggerExpressionUtils.java
@@ -0,0 +1,118 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.gravitino.maintenance.optimizer.recommender.util;
+
+import java.util.Map;
+import org.apache.gravitino.rel.Column;
+import org.apache.gravitino.rel.Table;
+import org.apache.gravitino.rel.expressions.sorts.SortOrder;
+import org.apache.gravitino.rel.expressions.transforms.Transform;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+class TestTableMetadataTriggerExpressionUtils {
+
+ @Test
+ void testBuildTableMetadataContextWithAllFields() {
+ Table tableMetadata = Mockito.mock(Table.class);
+ Mockito.when(tableMetadata.columns()).thenReturn(new Column[3]);
+ Mockito.when(tableMetadata.partitioning()).thenReturn(new Transform[1]);
+ Mockito.when(tableMetadata.sortOrder()).thenReturn(new SortOrder[1]);
+ Mockito.when(tableMetadata.properties())
+ .thenReturn(Map.of("format", "iceberg", "max_file_size",
"1073741824"));
+
+ Map<String, Object> context =
+
TableMetadataTriggerExpressionUtils.buildTableMetadataContext(tableMetadata);
+
+ Assertions.assertEquals(5, context.size());
+ Assertions.assertEquals(3, context.get("column_count"));
+ Assertions.assertEquals(1, context.get("partition_count"));
+ Assertions.assertEquals(1, context.get("sort_order_count"));
+ Assertions.assertEquals("iceberg", context.get("format"));
+ Assertions.assertEquals(1073741824L, context.get("max_file_size"));
+ }
+
+ @Test
+ void testBuildTableMetadataContextWithEmptyTableMetadata() {
+ Table tableMetadata = Mockito.mock(Table.class);
+ Mockito.when(tableMetadata.columns()).thenReturn(null);
+ Mockito.when(tableMetadata.partitioning()).thenReturn(null);
+ Mockito.when(tableMetadata.sortOrder()).thenReturn(null);
+ Mockito.when(tableMetadata.properties()).thenReturn(Map.of());
+
+ Map<String, Object> context =
+
TableMetadataTriggerExpressionUtils.buildTableMetadataContext(tableMetadata);
+
+ Assertions.assertEquals(3, context.size());
+ Assertions.assertEquals(0, context.get("column_count"));
+ Assertions.assertEquals(0, context.get("partition_count"));
+ Assertions.assertEquals(0, context.get("sort_order_count"));
+ }
+
+ @Test
+ void testBuildTableMetadataContextWithNullProperties() {
+ Table tableMetadata = Mockito.mock(Table.class);
+ Mockito.when(tableMetadata.columns()).thenReturn(new Column[1]);
+ Mockito.when(tableMetadata.partitioning()).thenReturn(new Transform[0]);
+ Mockito.when(tableMetadata.sortOrder()).thenReturn(new SortOrder[0]);
+ Mockito.when(tableMetadata.properties()).thenReturn(null);
+
+ Map<String, Object> context =
+
TableMetadataTriggerExpressionUtils.buildTableMetadataContext(tableMetadata);
+
+ Assertions.assertEquals(3, context.size());
+ Assertions.assertEquals(1, context.get("column_count"));
+ Assertions.assertEquals(0, context.get("partition_count"));
+ Assertions.assertEquals(0, context.get("sort_order_count"));
+ }
+
+ @Test
+ void testBuildTableMetadataContextParsesNumericProperties() {
+ Table tableMetadata = Mockito.mock(Table.class);
+ Mockito.when(tableMetadata.columns()).thenReturn(new Column[1]);
+ Mockito.when(tableMetadata.partitioning()).thenReturn(new Transform[0]);
+ Mockito.when(tableMetadata.sortOrder()).thenReturn(new SortOrder[0]);
+ Mockito.when(tableMetadata.properties())
+ .thenReturn(
+ Map.of(
+ "count", "12345",
+ "size", "9876543210",
+ "name", "test_table",
+ "enabled", "true"));
+
+ Map<String, Object> context =
+
TableMetadataTriggerExpressionUtils.buildTableMetadataContext(tableMetadata);
+
+ Assertions.assertEquals(12345L, context.get("count"));
+ Assertions.assertEquals(9876543210L, context.get("size"));
+ Assertions.assertEquals("test_table", context.get("name"));
+ Assertions.assertEquals("true", context.get("enabled"));
+ }
+
+ @Test
+ void testBuildTableMetadataContextFailsWhenTableMetadataIsNull() {
+ IllegalArgumentException exception =
+ Assertions.assertThrows(
+ IllegalArgumentException.class,
+ () ->
TableMetadataTriggerExpressionUtils.buildTableMetadataContext(null));
+ Assertions.assertTrue(exception.getMessage().contains("Table metadata is
null"));
+ }
+}