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

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


The following commit(s) were added to refs/heads/master by this push:
     new 20c5508add7 [FLINK-40022][table] Fix TableResult.await() hang on empty 
SELECT result
20c5508add7 is described below

commit 20c5508add753d55c7228e4679cd4e52612b5de1
Author: Timo Theusner <[email protected]>
AuthorDate: Thu Jul 30 10:00:44 2026 +0200

    [FLINK-40022][table] Fix TableResult.await() hang on empty SELECT result
    
    This closes #28585.
---
 .../flink/table/api/internal/ResultProvider.java   |  8 +-
 .../planner/connectors/CollectDynamicSink.java     | 11 ++-
 .../connectors/CollectResultProviderITCase.java    | 89 ++++++++++++++++++++++
 3 files changed, 102 insertions(+), 6 deletions(-)

diff --git 
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/internal/ResultProvider.java
 
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/internal/ResultProvider.java
index 061ef09d46d..1c23bed724d 100644
--- 
a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/internal/ResultProvider.java
+++ 
b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/internal/ResultProvider.java
@@ -53,10 +53,12 @@ public interface ResultProvider {
     RowDataToStringConverter getRowDataStringConverter();
 
     /**
-     * Return true if the first row is ready.
+     * Returns {@code true} once the result is ready to be consumed.
      *
-     * <p>The first row is ready when {@link CloseableIterator#hasNext} method 
returns true or
-     * {@link CloseableIterator#next()} method returns a row.
+     * <p>The result is ready when {@link CloseableIterator#hasNext()} returns 
{@code true} (a first
+     * row can be accessed) or {@link CloseableIterator#next()} returns a row, 
and when {@link
+     * CloseableIterator#hasNext()} returns {@code false} because the job has 
finished without
+     * producing rows.
      */
     boolean isFirstRowReady();
 
diff --git 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/connectors/CollectDynamicSink.java
 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/connectors/CollectDynamicSink.java
index 672dc7dfd7b..fda19f414c3 100644
--- 
a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/connectors/CollectDynamicSink.java
+++ 
b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/connectors/CollectDynamicSink.java
@@ -217,9 +217,14 @@ public final class CollectDynamicSink implements 
DynamicTableSink {
 
         @Override
         public boolean isFirstRowReady() {
-            return (this.rowDataIterator != null && 
this.rowDataIterator.firstRowProcessed)
-                    || (this.rowIterator != null && 
this.rowIterator.firstRowProcessed)
-                    || iterator.hasNext();
+            if ((this.rowDataIterator != null && 
this.rowDataIterator.firstRowProcessed)
+                    || (this.rowIterator != null && 
this.rowIterator.firstRowProcessed)) {
+                return true;
+            }
+            // hasNext() blocks until the first row is available or the job 
terminates. Once it
+            // returns we have a definitive answer, so the result is ready to 
be consumed.
+            iterator.hasNext();
+            return true;
         }
 
         @Override
diff --git 
a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/connectors/CollectResultProviderITCase.java
 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/connectors/CollectResultProviderITCase.java
new file mode 100644
index 00000000000..4ba9f0ff4c9
--- /dev/null
+++ 
b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/connectors/CollectResultProviderITCase.java
@@ -0,0 +1,89 @@
+/*
+ * 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.flink.table.planner.connectors;
+
+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.test.junit5.MiniClusterExtension;
+import org.apache.flink.types.Row;
+import org.apache.flink.util.CloseableIterator;
+import org.apache.flink.util.CollectionUtil;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
+import org.junit.jupiter.api.extension.ExtendWith;
+
+import java.util.concurrent.TimeUnit;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * ITCase for collecting SELECT results via the Table API (backed by {@code
+ * CollectDynamicSink.CollectResultProvider}).
+ */
+@ExtendWith(MiniClusterExtension.class)
+@Timeout(value = 60, unit = TimeUnit.SECONDS)
+class CollectResultProviderITCase {
+
+    private static final String EMPTY_RESULT_QUERY =
+            "SELECT * FROM (VALUES (1)) AS t(x) WHERE x < 0";
+
+    @Test
+    void awaitAndCollectCompleteForEmptyBatchResult() throws Exception {
+        final TableEnvironment tEnv =
+                
TableEnvironment.create(EnvironmentSettings.newInstance().inBatchMode().build());
+
+        final TableResult result = tEnv.executeSql(EMPTY_RESULT_QUERY);
+
+        result.await();
+        try (CloseableIterator<Row> rows = result.collect()) {
+            assertThat(rows.hasNext()).isFalse();
+        }
+    }
+
+    @Test
+    void awaitAndCollectCompleteForNonEmptyBatchResult() throws Exception {
+        final TableEnvironment tEnv =
+                
TableEnvironment.create(EnvironmentSettings.newInstance().inBatchMode().build());
+
+        final TableResult result =
+                tEnv.executeSql("SELECT x FROM (VALUES (1), (2), (3)) AS t(x) 
WHERE x > 1");
+
+        result.await();
+        try (CloseableIterator<Row> rows = result.collect()) {
+            assertThat(CollectionUtil.iteratorToList(rows))
+                    .containsExactlyInAnyOrder(Row.of(2), Row.of(3));
+        }
+    }
+
+    @Test
+    void awaitAndCollectCompleteForEmptyStreamingResult() throws Exception {
+        final TableEnvironment tEnv =
+                TableEnvironment.create(
+                        
EnvironmentSettings.newInstance().inStreamingMode().build());
+
+        final TableResult result = tEnv.executeSql(EMPTY_RESULT_QUERY);
+
+        result.await();
+        try (CloseableIterator<Row> rows = result.collect()) {
+            assertThat(rows.hasNext()).isFalse();
+        }
+    }
+}

Reply via email to