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();
+ }
+ }
+}