This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 5652c0023c [test] Fix unstable unaware bucket read test and assert 
split distribution (#8850)
5652c0023c is described below

commit 5652c0023ce7455553eedd721011054bd8213d3a
Author: Vova Kolmakov <[email protected]>
AuthorDate: Mon Jul 27 15:05:38 2026 +0700

    [test] Fix unstable unaware bucket read test and assert split distribution 
(#8850)
---
 .../apache/paimon/flink/AppendOnlyTableITCase.java |  16 ---
 .../source/UnawareBucketSplitShuffleITCase.java    | 147 +++++++++++++++++++++
 2 files changed, 147 insertions(+), 16 deletions(-)

diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/AppendOnlyTableITCase.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/AppendOnlyTableITCase.java
index f14c9854df..d9e0987fa6 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/AppendOnlyTableITCase.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/AppendOnlyTableITCase.java
@@ -405,22 +405,6 @@ public class AppendOnlyTableITCase extends 
CatalogITCaseBase {
                 .isGreaterThanOrEqualTo(1);
     }
 
-    @Test
-    public void testReadUnawareBucketTableWithRebalanceShuffle() throws 
Exception {
-        batchSql(
-                "CREATE TABLE append_scalable_table (id INT, data STRING) "
-                        + "WITH ('bucket' = '-1', 'consumer-id' = 'test', 
'consumer.expiration-time' = '365 d', 'target-file-size' = '1 B', 
'source.split.target-size' = '1 B', 'scan.parallelism' = '4')");
-        batchSql("INSERT INTO append_scalable_table VALUES (1, 'AAA'), (2, 
'BBB')");
-        batchSql("INSERT INTO append_scalable_table VALUES (1, 'AAA'), (2, 
'BBB')");
-        batchSql("INSERT INTO append_scalable_table VALUES (1, 'AAA'), (2, 
'BBB')");
-        batchSql("INSERT INTO append_scalable_table VALUES (1, 'AAA'), (2, 
'BBB')");
-
-        BlockingIterator<Row, Row> iterator =
-                BlockingIterator.of(streamSqlIter(("SELECT id FROM 
append_scalable_table")));
-        assertThat(iterator.collect(2)).containsExactlyInAnyOrder(Row.of(1), 
Row.of(2));
-        iterator.close();
-    }
-
     @Test
     public void testReadPartitionOrder() {
         setParallelism(1);
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/UnawareBucketSplitShuffleITCase.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/UnawareBucketSplitShuffleITCase.java
new file mode 100644
index 0000000000..d71a213fbe
--- /dev/null
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/source/UnawareBucketSplitShuffleITCase.java
@@ -0,0 +1,147 @@
+/*
+ * 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.paimon.flink.source;
+
+import org.apache.paimon.flink.util.MiniClusterWithClientExtension;
+import org.apache.paimon.utils.BlockingIterator;
+
+import org.apache.flink.api.common.JobID;
+import org.apache.flink.client.program.ClusterClient;
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.metrics.groups.OperatorMetricGroup;
+import org.apache.flink.runtime.client.JobStatusMessage;
+import org.apache.flink.runtime.testutils.InMemoryReporter;
+import org.apache.flink.runtime.testutils.MiniClusterResourceConfiguration;
+import org.apache.flink.table.api.EnvironmentSettings;
+import org.apache.flink.table.api.TableEnvironment;
+import org.apache.flink.table.api.TableResult;
+import org.apache.flink.table.api.config.ExecutionConfigOptions;
+import org.apache.flink.types.Row;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+import org.junit.jupiter.api.io.TempDir;
+
+import java.nio.file.Path;
+import java.util.List;
+import java.util.stream.Collectors;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * IT case for the split shuffle of a streaming read over an unaware bucket 
table. Such a table has
+ * a single bucket, so a bucket based shuffle would send every split to one 
reader; splits must be
+ * rebalanced over all of them instead.
+ */
+public class UnawareBucketSplitShuffleITCase {
+
+    private static final int READER_PARALLELISM = 4;
+
+    // matches the reader operator of every subtask, the planner names it 
"table_name[id]"
+    private static final String READER_METRIC_PATTERN = 
"append_scalable_table\\[";
+
+    private static final InMemoryReporter reporter = 
InMemoryReporter.createWithRetainedMetrics();
+
+    @RegisterExtension
+    protected static final MiniClusterWithClientExtension 
MINI_CLUSTER_EXTENSION =
+            new MiniClusterWithClientExtension(
+                    new MiniClusterResourceConfiguration.Builder()
+                            .setNumberTaskManagers(1)
+                            .setNumberSlotsPerTaskManager(READER_PARALLELISM)
+                            .setConfiguration(reporter.addToConfiguration(new 
Configuration()))
+                            .build());
+
+    @TempDir Path tempPath;
+
+    @AfterEach
+    public final void cleanupRunningJobs() throws Exception {
+        ClusterClient<?> clusterClient = 
MINI_CLUSTER_EXTENSION.createRestClusterClient();
+        for (JobStatusMessage job : clusterClient.listJobs().get()) {
+            if (!job.getJobState().isTerminalState()) {
+                try {
+                    clusterClient.cancel(job.getJobId()).get();
+                } catch (Exception ignored) {
+                    // ignore exceptions when cancelling dangling jobs
+                }
+            }
+        }
+    }
+
+    @Test
+    public void testReadUnawareBucketTableWithRebalanceShuffle() throws 
Exception {
+        TableEnvironment batchEnv = tableEnv(false);
+        batchEnv.executeSql(
+                "CREATE TABLE append_scalable_table (id INT, data STRING) "
+                        + "WITH ('bucket' = '-1', 'consumer-id' = 'test', 
'consumer.expiration-time' = '365 d', "
+                        + "'target-file-size' = '1 B', 
'source.split.target-size' = '1 B', 'scan.parallelism' = '4')");
+        for (int i = 0; i < READER_PARALLELISM; i++) {
+            batchEnv.executeSql("INSERT INTO append_scalable_table VALUES (1, 
'AAA'), (2, 'BBB')")
+                    .await();
+        }
+
+        TableResult result = tableEnv(true).executeSql("SELECT id FROM 
append_scalable_table");
+        JobID jobId = result.getJobClient().get().getJobID();
+
+        try (BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(result.collect())) {
+            // collect every row, which record arrives first is not 
deterministic once the splits
+            // are spread over concurrent readers
+            assertThat(iterator.collect(2 * READER_PARALLELISM))
+                    .containsExactlyInAnyOrder(
+                            Row.of(1), Row.of(1), Row.of(1), Row.of(1), 
Row.of(2), Row.of(2),
+                            Row.of(2), Row.of(2));
+
+            // the property under test: all rows are in, so every split has 
been handed out, and
+            // the rebalance must have left no reader empty. The bucket 
shuffle used before #3955
+            // sends every split of a single bucket table to one channel and 
gives [0, 0, 0, 8]
+            assertThat(recordsReadPerReaderSubtask(jobId))
+                    .hasSize(READER_PARALLELISM)
+                    .allSatisfy(records -> assertThat(records).isPositive());
+        }
+    }
+
+    // the reporter registers one metric group per reader subtask, so one 
entry per subtask
+    private List<Long> recordsReadPerReaderSubtask(JobID jobId) throws 
Exception {
+        for (int tries = 1; tries <= 20; tries++) {
+            List<OperatorMetricGroup> groups =
+                    reporter.findOperatorMetricGroups(jobId, 
READER_METRIC_PATTERN);
+            if (groups.size() == READER_PARALLELISM) {
+                return groups.stream()
+                        .map(group -> 
group.getIOMetricGroup().getNumRecordsInCounter().getCount())
+                        .collect(Collectors.toList());
+            }
+            Thread.sleep(1000);
+        }
+        throw new AssertionError(
+                "Reader metric groups of job " + jobId + " did not show up, 
pattern may be stale");
+    }
+
+    private TableEnvironment tableEnv(boolean streaming) {
+        EnvironmentSettings.Builder settings = 
EnvironmentSettings.newInstance();
+        TableEnvironment tEnv =
+                TableEnvironment.create(
+                        (streaming ? settings.inStreamingMode() : 
settings.inBatchMode()).build());
+        
tEnv.getConfig().set(ExecutionConfigOptions.TABLE_EXEC_RESOURCE_DEFAULT_PARALLELISM,
 2);
+        tEnv.executeSql(
+                "CREATE CATALOG mycat WITH ( 'type' = 'paimon', 'warehouse' = 
'"
+                        + tempPath
+                        + "' )");
+        tEnv.useCatalog("mycat");
+        return tEnv;
+    }
+}

Reply via email to