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


##########
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;
+        }
+    }
+
+    /**
+     * Waits until the last stable offset of {@code tp} reaches {@code 
expectedOffset}, and returns
+     * it. In most cases, this will return the expected value on the first 
check, but the
+     * transaction coordinator propagates it to partition leaders 
asynchronously, so LSO can briefly
+     * lag behind a committed transaction. To avoid introducing a test race 
condition, this method
+     * checks again after a brief wait.
+     */
+    private static long awaitLastStableOffset(TopicPartition tp, long 
expectedOffset)
+            throws Exception {
+        try (KafkaConsumer<String, String> probe = createReadCommittedProbe()) 
{
+            final List<TopicPartition> partitions = 
Collections.singletonList(tp);
+            long lastStableOffset = probe.endOffsets(partitions).get(tp);
+            final long deadline = System.currentTimeMillis() + 30_000L;
+            while (lastStableOffset != expectedOffset && 
System.currentTimeMillis() < deadline) {
+                Thread.sleep(200L);

Review Comment:
   Agreed - updated



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