dalelane commented on code in PR #293:
URL:
https://github.com/apache/flink-connector-kafka/pull/293#discussion_r4161789230
##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/reader/KafkaSourceReaderTest.java:
##########
@@ -695,6 +1169,72 @@ private long getCommittedOffsetMetric(TopicPartition tp,
MetricListener listener
// ---------------------
+ private static KafkaConsumer<String, String> createReadCommittedProbe() {
+ final Properties props = new Properties();
+ props.putAll(KafkaSourceTestEnv.standardProps);
+ props.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "read-probe-" +
UUID.randomUUID());
+ props.setProperty(ConsumerConfig.ISOLATION_LEVEL_CONFIG,
"read_committed");
+ props.setProperty(
+ ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName());
+ props.setProperty(
+ ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName());
+ return new KafkaConsumer<>(props);
+ }
+
+ private static int getReadCommittedRecordsCount(TopicPartition tp) {
+ try (KafkaConsumer<String, String> probe = createReadCommittedProbe())
{
+ final List<TopicPartition> partitions =
Collections.singletonList(tp);
+ probe.assign(partitions);
+ probe.seekToBeginning(partitions);
+ final long lastStableOffset = probe.endOffsets(partitions).get(tp);
+ int count = 0;
+ final long deadline = System.currentTimeMillis() + 30_000L;
+ while (probe.position(tp) < lastStableOffset &&
System.currentTimeMillis() < deadline) {
+ List<ConsumerRecord<String, String>> records =
+ probe.poll(Duration.ofMillis(500)).records(tp);
+ count += records.size();
+ }
+ return count;
Review Comment:
Good point. It now uses waitUtil as well, so a timeout fails with "The probe
did not reach the last stable offset"
--
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]