This is an automated email from the ASF dual-hosted git repository.
voonhous pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git
The following commit(s) were added to refs/heads/master by this push:
new 8a41388d76a8 fix(trino): keep none() when splitting predicates (#19863)
8a41388d76a8 is described below
commit 8a41388d76a89f8316dfbf787b8f3188f717a8f4
Author: voonhous <[email protected]>
AuthorDate: Tue Sep 8 13:38:05 2026 +0800
fix(trino): keep none() when splitting predicates (#19863)
Port of trinodb/trino#30958 (03074fb642d9). HudiPredicates.from only
read TupleDomain.getDomains(), which is empty for both none() and all(),
so a none() constraint came back as all() on both the partition and the
regular side. Return none() for both when the input is none(), and add
upstream's TestHudiPredicates.
One addition over upstream: HudiSplitManager.getSplits now returns an
empty split source when either predicate on the table handle is none().
Without it, the preserved none() would reach MetastoreUtil's
computePartitionKeyFilter, whose checkArgument rejects a none() domain,
turning a query that should produce no splits into a failure. The
planner already turns a none() domain into an empty ValuesNode before
applyFilter, so neither path is reachable through the engine today; the
guard is the same defense HiveSplitManager carries.
---
.../java/io/trino/plugin/hudi/HudiPredicates.java | 4 +
.../io/trino/plugin/hudi/HudiSplitManager.java | 7 ++
.../io/trino/plugin/hudi/TestHudiPredicates.java | 98 ++++++++++++++++++++++
.../io/trino/plugin/hudi/TestHudiSplitSource.java | 50 +++++++++++
4 files changed, 159 insertions(+)
diff --git a/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiPredicates.java
b/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiPredicates.java
index 16859e8e6dd3..5161a85ea1b6 100644
--- a/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiPredicates.java
+++ b/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiPredicates.java
@@ -29,6 +29,10 @@ public class HudiPredicates
public static HudiPredicates from(TupleDomain<ColumnHandle> predicate)
{
+ if (predicate.isNone()) {
+ return new HudiPredicates(TupleDomain.none(), TupleDomain.none());
+ }
+
Map<HiveColumnHandle, Domain> partitionColumnPredicates = new
HashMap<>();
Map<HiveColumnHandle, Domain> regularColumnPredicates = new
HashMap<>();
diff --git
a/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiSplitManager.java
b/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiSplitManager.java
index df0018a49268..223894f4e2a5 100644
--- a/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiSplitManager.java
+++ b/hudi-trino/src/main/java/io/trino/plugin/hudi/HudiSplitManager.java
@@ -51,6 +51,7 @@ import static
io.trino.plugin.hudi.HudiSessionProperties.getDynamicFilteringWait
import static
io.trino.plugin.hudi.HudiSessionProperties.getMaxOutstandingSplits;
import static io.trino.plugin.hudi.HudiSessionProperties.getMaxSplitsPerSecond;
import static
io.trino.plugin.hudi.partition.HiveHudiPartitionInfo.NON_PARTITION;
+import static io.trino.spi.connector.FixedSplitSource.emptySplitSource;
import static java.lang.String.format;
import static java.util.Objects.requireNonNull;
@@ -82,6 +83,12 @@ public class HudiSplitManager
Constraint constraint)
{
HudiTableHandle hudiTableHandle = (HudiTableHandle) tableHandle;
+ // The planner turns a none() domain into an empty ValuesNode before
applyFilter, so a none() predicate
+ // should never reach here; guard anyway, as computePartitionKeyFilter
below throws on a none() domain.
+ if (hudiTableHandle.getPartitionPredicates().isNone() ||
hudiTableHandle.getRegularPredicates().isNone()) {
+ return emptySplitSource();
+ }
+
HiveMetastore metastore =
metastoreProvider.apply(session.getIdentity(), (HiveTransactionHandle)
transaction);
Lazy<Map<String, Partition>> lazyAllPartitions = Lazy.lazily(() -> {
HoodieTimer timer = HoodieTimer.start();
diff --git
a/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiPredicates.java
b/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiPredicates.java
new file mode 100644
index 000000000000..a7b74f2b958b
--- /dev/null
+++ b/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiPredicates.java
@@ -0,0 +1,98 @@
+/*
+ * Licensed 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 io.trino.plugin.hudi;
+
+import com.google.common.collect.ImmutableMap;
+import io.airlift.slice.Slices;
+import io.trino.plugin.hive.HiveColumnHandle;
+import io.trino.spi.connector.ColumnHandle;
+import io.trino.spi.predicate.Domain;
+import io.trino.spi.predicate.TupleDomain;
+import org.junit.jupiter.api.Test;
+
+import java.util.Optional;
+
+import static io.trino.metastore.HiveType.HIVE_STRING;
+import static io.trino.plugin.hive.HiveColumnHandle.ColumnType.PARTITION_KEY;
+import static io.trino.plugin.hive.HiveColumnHandle.ColumnType.REGULAR;
+import static io.trino.plugin.hive.HiveColumnHandle.createBaseColumn;
+import static io.trino.spi.type.VarcharType.VARCHAR;
+import static org.assertj.core.api.Assertions.assertThat;
+
+public class TestHudiPredicates
+{
+ private static final HiveColumnHandle PARTITION_COLUMN =
createBaseColumn("part", 0, HIVE_STRING, VARCHAR, PARTITION_KEY,
Optional.empty());
+ private static final HiveColumnHandle REGULAR_COLUMN =
createBaseColumn("col", 1, HIVE_STRING, VARCHAR, REGULAR, Optional.empty());
+
+ @Test
+ public void testNoneConstraintStaysNone()
+ {
+ HudiPredicates predicates = HudiPredicates.from(TupleDomain.none());
+
+
assertThat(predicates.getPartitionColumnPredicates().isNone()).isTrue();
+ assertThat(predicates.getRegularColumnPredicates().isNone()).isTrue();
+ }
+
+ @Test
+ public void testAllConstraintStaysAll()
+ {
+ HudiPredicates predicates = HudiPredicates.from(TupleDomain.all());
+
+ assertThat(predicates.getPartitionColumnPredicates().isAll()).isTrue();
+ assertThat(predicates.getRegularColumnPredicates().isAll()).isTrue();
+ }
+
+ @Test
+ public void testConstraintSplitByColumnType()
+ {
+ Domain partitionDomain = Domain.singleValue(VARCHAR,
Slices.utf8Slice("p1"));
+ Domain regularDomain = Domain.singleValue(VARCHAR,
Slices.utf8Slice("v1"));
+ TupleDomain<ColumnHandle> predicate =
TupleDomain.withColumnDomains(ImmutableMap.of(
+ PARTITION_COLUMN, partitionDomain,
+ REGULAR_COLUMN, regularDomain));
+
+ HudiPredicates predicates = HudiPredicates.from(predicate);
+
+ assertThat(predicates.getPartitionColumnPredicates())
+
.isEqualTo(TupleDomain.withColumnDomains(ImmutableMap.of(PARTITION_COLUMN,
partitionDomain)));
+ assertThat(predicates.getRegularColumnPredicates())
+
.isEqualTo(TupleDomain.withColumnDomains(ImmutableMap.of(REGULAR_COLUMN,
regularDomain)));
+ }
+
+ @Test
+ public void testOnlyPartitionColumnConstraintLeavesRegularAll()
+ {
+ Domain partitionDomain = Domain.singleValue(VARCHAR,
Slices.utf8Slice("p1"));
+ TupleDomain<ColumnHandle> predicate =
TupleDomain.withColumnDomains(ImmutableMap.of(PARTITION_COLUMN,
partitionDomain));
+
+ HudiPredicates predicates = HudiPredicates.from(predicate);
+
+ assertThat(predicates.getPartitionColumnPredicates())
+
.isEqualTo(TupleDomain.withColumnDomains(ImmutableMap.of(PARTITION_COLUMN,
partitionDomain)));
+ assertThat(predicates.getRegularColumnPredicates().isAll()).isTrue();
+ }
+
+ @Test
+ public void testOnlyRegularColumnConstraintLeavesPartitionAll()
+ {
+ Domain regularDomain = Domain.singleValue(VARCHAR,
Slices.utf8Slice("v1"));
+ TupleDomain<ColumnHandle> predicate =
TupleDomain.withColumnDomains(ImmutableMap.of(REGULAR_COLUMN, regularDomain));
+
+ HudiPredicates predicates = HudiPredicates.from(predicate);
+
+ assertThat(predicates.getPartitionColumnPredicates().isAll()).isTrue();
+ assertThat(predicates.getRegularColumnPredicates())
+
.isEqualTo(TupleDomain.withColumnDomains(ImmutableMap.of(REGULAR_COLUMN,
regularDomain)));
+ }
+}
diff --git
a/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiSplitSource.java
b/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiSplitSource.java
index 567d4b776459..989e8ebc6a62 100644
--- a/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiSplitSource.java
+++ b/hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiSplitSource.java
@@ -13,24 +13,33 @@
*/
package io.trino.plugin.hudi;
+import com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableSet;
import io.airlift.units.Duration;
+import io.trino.plugin.hive.HiveColumnHandle;
+import io.trino.plugin.hive.HiveTransactionHandle;
import io.trino.plugin.hive.util.AsyncQueue;
import io.trino.plugin.hive.util.ThrottledAsyncQueue;
import io.trino.spi.TrinoException;
import io.trino.spi.connector.ConnectorSplit;
+import io.trino.spi.connector.ConnectorSplitSource;
+import io.trino.spi.connector.Constraint;
import io.trino.spi.connector.DynamicFilterSnapshot;
import io.trino.spi.connector.SchemaTableName;
import io.trino.spi.predicate.TupleDomain;
+import org.apache.hudi.common.model.HoodieTableType;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.TestInstance;
import java.util.List;
+import java.util.OptionalLong;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
+import static io.trino.testing.TestingConnectorSession.SESSION;
import static java.util.concurrent.TimeUnit.SECONDS;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -68,6 +77,29 @@ public class TestHudiSplitSource
assertThat(splitSource.isFinished()).isTrue();
}
+ @Test
+ public void testNonePredicateYieldsEmptySplitSource()
+ {
+ // A none() predicate must short-circuit before the metastore is
consulted: the split manager would
+ // otherwise hand it to computePartitionKeyFilter, which rejects a
none() domain
+ HudiSplitManager splitManager = new HudiSplitManager(
+ (identity, transaction) -> {
+ throw new AssertionError("metastore must not be consulted
for a none() predicate");
+ },
+ executor,
+ scheduler);
+
+ for (HudiTableHandle tableHandle : List.of(
+ createTableHandle(TupleDomain.none(), TupleDomain.all()),
+ createTableHandle(TupleDomain.all(), TupleDomain.none()))) {
+ ConnectorSplitSource splitSource = splitManager.getSplits(
+ new HiveTransactionHandle(false), SESSION, tableHandle,
ImmutableSet.of(), Constraint.alwaysTrue());
+
+ assertThat(splitSource.isFinished()).isTrue();
+ assertThat(splitSource.getNextBatch(10, new
DynamicFilterSnapshot(TupleDomain.all(), false)).join()).isEmpty();
+ }
+ }
+
@Test
public void testLoaderFailureBlocksFinishedAndSurfaces()
throws Exception
@@ -118,4 +150,22 @@ public class TestHudiSplitSource
}
assertThat(queue.isFinished()).isTrue();
}
+
+ private static HudiTableHandle createTableHandle(
+ TupleDomain<HiveColumnHandle> partitionPredicates,
+ TupleDomain<HiveColumnHandle> regularPredicates)
+ {
+ return new HudiTableHandle(
+ TABLE.getSchemaName(),
+ TABLE.getTableName(),
+ "/test/path",
+ HoodieTableType.COPY_ON_WRITE,
+ ImmutableList.of(),
+ ImmutableList.of(),
+ partitionPredicates,
+ regularPredicates,
+ OptionalLong.empty(),
+ "",
+ "101");
+ }
}