From c4fb691b3fa89a9ef678eafedc81f6e05045ec5d Mon Sep 17 00:00:00 2001 From: WenDing-Y <1062698930@qq.com> Date: Tue, 11 Aug 2026 19:42:51 +0800 Subject: [PATCH] [flink] Fix BETWEEN literal extraction in PredicateConverter --- .../fluss/flink/utils/PredicateConverter.java | 6 +++- .../flink/utils/PredicateConverterTest.java | 29 ++++++++++++++++++- 2 files changed, 33 insertions(+), 2 deletions(-) diff --git a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/PredicateConverter.java b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/PredicateConverter.java index 7c691cea978..1b75f210878 100644 --- a/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/PredicateConverter.java +++ b/fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/utils/PredicateConverter.java @@ -130,8 +130,12 @@ public Predicate visit(CallExpression call) { } else if (func == BuiltInFunctionDefinitions.BETWEEN) { FieldReferenceExpression fieldRefExpr = extractFieldReference(children.get(0)).orElseThrow(UnsupportedExpression::new); + Object lowerBound = + extractLiteral(fieldRefExpr.getOutputDataType(), children.get(1)); + Object upperBound = + extractLiteral(fieldRefExpr.getOutputDataType(), children.get(2)); return builder.between( - builder.indexOf(fieldRefExpr.getName()), children.get(1), children.get(2)); + builder.indexOf(fieldRefExpr.getName()), lowerBound, upperBound); } else if (func == BuiltInFunctionDefinitions.LIKE) { FieldReferenceExpression fieldRefExpr = extractFieldReference(children.get(0)).orElseThrow(UnsupportedExpression::new); diff --git a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/utils/PredicateConverterTest.java b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/utils/PredicateConverterTest.java index 76c7f5d68b6..85506959c76 100644 --- a/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/utils/PredicateConverterTest.java +++ b/fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/utils/PredicateConverterTest.java @@ -72,6 +72,33 @@ public class PredicateConverterTest { private static final PredicateConverter CONVERTER = new PredicateConverter(BUILDER); + /** + * Stronger than toString equality: BETWEEN must evaluate against rows. The parameterized + * toString check can pass even when Flink {@code ValueLiteralExpression} nodes were stored as + * literals without {@code extractLiteral}. + */ + @Test + public void testBetweenEvaluatesAgainstRow() { + FieldReferenceExpression longRefExpr = + new FieldReferenceExpression( + "long1", DataTypes.BIGINT(), Integer.MAX_VALUE, Integer.MAX_VALUE); + CallExpression between = + CallExpression.permanent( + BuiltInFunctionDefinitions.BETWEEN, + Arrays.asList( + longRefExpr, + new ValueLiteralExpression(10), + new ValueLiteralExpression(20)), + DataTypes.BOOLEAN()); + + Predicate predicate = between.accept(CONVERTER); + assertThat(predicate.test(GenericRow.of(15L))).isTrue(); + assertThat(predicate.test(GenericRow.of(10L))).isTrue(); + assertThat(predicate.test(GenericRow.of(20L))).isTrue(); + assertThat(predicate.test(GenericRow.of(9L))).isFalse(); + assertThat(predicate.test(GenericRow.of(21L))).isFalse(); + } + @MethodSource("provideResolvedExpression") @ParameterizedTest public void testVisitAndAutoTypeInference(ResolvedExpression expression, Predicate expected) { @@ -252,7 +279,7 @@ public static Stream provideResolvedExpression() { BuiltInFunctionDefinitions.BETWEEN, Arrays.asList(longRefExpr, intLitExpr, intLitExpr2), DataTypes.BOOLEAN()), - BUILDER.between(0, 10, 20))); + BUILDER.between(0, 10L, 20L))); } @MethodSource("provideLikeExpressions")