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

czweng pushed a commit to branch release-1.17
in repository https://gitbox.apache.org/repos/asf/flink.git


The following commit(s) were added to refs/heads/release-1.17 by this push:
     new bb66b989188 [FLINK-25813][table-planner] Fix unstable test 
TableITCase#testCollectWithClose
bb66b989188 is described below

commit bb66b989188c661d9d7847b6a5728697a3a2dcd2
Author: Matthias Pohl <[email protected]>
AuthorDate: Tue Feb 14 03:26:19 2023 +0100

    [FLINK-25813][table-planner] Fix unstable test 
TableITCase#testCollectWithClose
    
    This closes #21921.
---
 .../org/apache/flink/table/api/TableITCase.scala   | 37 ++++++++++++----------
 1 file changed, 21 insertions(+), 16 deletions(-)

diff --git 
a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/api/TableITCase.scala
 
b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/api/TableITCase.scala
index e108c98e265..e09c6274aab 100644
--- 
a/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/api/TableITCase.scala
+++ 
b/flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/api/TableITCase.scala
@@ -24,6 +24,7 @@ import 
org.apache.flink.table.api.internal.TableEnvironmentImpl
 import org.apache.flink.table.catalog.{Column, ResolvedSchema}
 import org.apache.flink.table.planner.utils.TestTableSourceSinks
 import org.apache.flink.test.util.AbstractTestBase
+import org.apache.flink.test.util.AbstractTestBase.MINI_CLUSTER_RESOURCE
 import org.apache.flink.types.{Row, RowKind}
 import org.apache.flink.util.CollectionUtil
 
@@ -32,7 +33,7 @@ import _root_.java.util
 import org.hamcrest.MatcherAssert.assertThat
 import org.hamcrest.Matchers.containsInAnyOrder
 import org.junit.{Before, Rule, Test}
-import org.junit.Assert.{assertEquals, assertNotEquals, assertTrue}
+import org.junit.Assert.{assertEquals, assertTrue}
 import org.junit.rules.TemporaryFolder
 import org.junit.runner.RunWith
 import org.junit.runners.Parameterized
@@ -108,28 +109,32 @@ class TableITCase(tableEnvName: String, isStreaming: 
Boolean) extends AbstractTe
 
   @Test
   def testCollectWithClose(): Unit = {
-    val query =
+    val sourceDdl =
       """
-        |select id, concat(concat(`first`, ' '), `last`) as `full name`
-        |from MyTable where mod(id, 2) = 0
-      """.stripMargin
+        |create table unbounded_source (
+        |  id int
+        |) with (
+        |  'connector' = 'datagen',
+        |  'number-of-rows' = '10000',
+        |  'rows-per-second' = '1' -- slow producing speed to make sure that
+        |                          -- source is not finished when job is 
cancelled
+        |)
+        |""".stripMargin
+    tEnv.executeSql(sourceDdl)
+    val query = "select id from unbounded_source where mod(id, 2) = 0"
     val table = tEnv.sqlQuery(query)
     val tableResult = table.execute()
     assertTrue(tableResult.getJobClient.isPresent)
     assertEquals(ResultKind.SUCCESS_WITH_CONTENT, tableResult.getResultKind)
     val it = tableResult.collect()
     it.close()
-    val jobStatus =
-      try {
-        Some(tableResult.getJobClient.get().getJobStatus.get())
-      } catch {
-        // ignore the exception,
-        // because the MiniCluster maybe already been shut down when getting 
job status
-        case _: Throwable => None
-      }
-    if (jobStatus.isDefined) {
-      assertNotEquals(JobStatus.RUNNING, jobStatus.get)
-    }
+
+    // wait for mini cluster to shut down
+    val jobClient = tableResult.getJobClient.get()
+    val jobId = jobClient.getJobID
+    MINI_CLUSTER_RESOURCE.getClusterClient.requestJobResult(jobId).get()
+
+    assertEquals(JobStatus.CANCELED, jobClient.getJobStatus.get())
   }
 
   @Test

Reply via email to