dalelane commented on code in PR #293:
URL: 
https://github.com/apache/flink-connector-kafka/pull/293#discussion_r4123969451


##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/reader/KafkaSourceReaderTest.java:
##########
@@ -266,6 +277,456 @@ void testOffsetCommitOnCheckpointComplete() throws 
Exception {
         }
     }
 
+    /** Writes the records that a {@link #offsetConvergenceScenarios} scenario 
needs. */
+    @FunctionalInterface
+    private interface RecordProducer {
+        void produce(String topic) throws Throwable;
+    }
+
+    private static Stream<Arguments> offsetConvergenceScenarios() {

Review Comment:
   Applied in a94226763cf24fca816c095b503ec2679f417b7f
   
   On its own that made `testOffsetCommitOnCheckpointComplete` and 
`testKafkaSourceMetrics` fail every time in my resource-constrained VM. Same 
sort of problem as above (tests that stop polling the reader before waiting for 
the commit and get blocked). I've added a `reader.pollNext(output)` to both of 
their wait loops, the same as the new tests already do.
   
   1 - 
https://github.com/apache/flink-connector-kafka/pull/293/changes#diff-912e48531349dc6ecdd8a8a2c25c5b65d4f3511bdd028a0e1cdb5605350759f3R254
   2 - 
https://github.com/apache/flink-connector-kafka/pull/293/changes#diff-912e48531349dc6ecdd8a8a2c25c5b65d4f3511bdd028a0e1cdb5605350759f3R829
   



-- 
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