gustavodemorais commented on code in PR #29296:
URL: https://github.com/apache/flink/pull/29296#discussion_r4132908707


##########
flink-table/flink-table-planner/src/test/scala/org/apache/flink/table/planner/runtime/stream/sql/JoinITCase.scala:
##########
@@ -1747,4 +1829,16 @@ object JoinITCase {
       Array(MiniBatchOff, HEAP_BACKEND, Boolean.box(true))
     )
   }
+
+  /**
+   * Emits the given rows once the sink received its first row, so that the 
other input goes first.
+   */
+  class WaitForSinkThenEmit(rows: Seq[Row]) extends 
GeneratorFunction[java.lang.Long, Row] {
+    override def map(index: java.lang.Long): Row = {
+      while (TestValuesTableFactory.getRawResultsAsStrings("sink").isEmpty) {

Review Comment:
   The test now generates the sink name and passes it to WaitForSinkThenEmit, 
which also removes the clearAllData() call. TestValues keys sink results by 
table name, so there's no sink id like registerData has



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to