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]