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


##########
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/reader/KafkaSourceReaderTest.java:
##########
@@ -266,6 +278,464 @@ 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() {
+        return Stream.of(
+                // no transaction and just one record
+                // offset 0 = record
+                // offset 1 << log end offset
+                Arguments.of(
+                        "NonTransactional",
+                        1L,
+                        1,
+                        1,
+                        (RecordProducer)
+                                topic ->
+                                        KafkaSourceTestEnv.produceToKafka(
+                                                Collections.singletonList(
+                                                        new ProducerRecord<>(
+                                                                topic, 0, 
topic + "-key", 0)))),
+                // transaction containing one record and a Kafka commit marker
+                // offset 0 = record
+                // offset 1 = commit marker
+                // offset 2 << log end offset
+                Arguments.of(
+                        "Transactional",
+                        2L,
+                        1,
+                        1,
+                        (RecordProducer)
+                                topic ->
+                                        KafkaSourceTestEnv.produceToKafka(
+                                                Collections.singletonList(
+                                                        new ProducerRecord<>(
+                                                                topic, 0, 
topic + "-key", 0)),
+                                                
kafkaServer.getTransactionalProducerConfig())),
+                // transaction containing three records
+                // offset 0 = record
+                // offset 1 = record
+                // offset 2 = record
+                // offset 3 = commit marker
+                // offset 4 << log end offset
+                Arguments.of(
+                        "MultiRecordTransaction",
+                        4L,
+                        3,
+                        1,
+                        (RecordProducer) topic -> produceInTransaction(topic, 
0, 3, true)),
+                // committed transaction containing two records
+                // aborted transaction containing three records
+                // offset 0 = record
+                // offset 1 = record
+                // offset 2 = commit marker
+                // offset 3 = record
+                // offset 4 = record
+                // offset 5 = record
+                // offset 6 = abort marker
+                // offset 7 << log end offset
+                Arguments.of(
+                        "AbortedTransaction",
+                        7L,
+                        2,
+                        1,
+                        (RecordProducer)
+                                topic -> {
+                                    produceInTransaction(topic, 0, 2, true);
+                                    produceInTransaction(topic, 0, 3, false);
+                                }),
+                // three separate transactions each containing one record
+                // offset 0 = record
+                // offset 1 = commit marker
+                // offset 2 = record
+                // offset 3 = commit marker
+                // offset 4 = record
+                // offset 5 = commit marker
+                // offset 6 << log end offset
+                Arguments.of(
+                        "MultipleTransactions",
+                        6L,
+                        3,
+                        1,
+                        (RecordProducer)
+                                topic -> {
+                                    produceInTransaction(topic, 0, 1, true);
+                                    produceInTransaction(topic, 0, 1, true);
+                                    produceInTransaction(topic, 0, 1, true);
+                                }),
+                // two transactional producers interleaving records on the 
same partition before
+                // either commits, so the two commit markers do not 
immediately follow their own
+                // records
+                // offset 0 = producer A record
+                // offset 1 = producer B record
+                // offset 2 = producer A record
+                // offset 3 = producer B record
+                // offset 4 = producer A commit marker
+                // offset 5 = producer B commit marker
+                // offset 6 << log end offset
+                Arguments.of(
+                        "InterleavedTransactions",
+                        6L,
+                        4,
+                        1,
+                        (RecordProducer) topic -> 
produceInterleavedTransactions(topic, 0, 2)),
+                // a single aborted transaction
+                // so the partition never delivers a record to the reader
+                // offset 0 = record
+                // offset 1 = record
+                // offset 2 = record
+                // offset 3 = abort marker
+                // offset 4 << log end offset
+                Arguments.of(
+                        "OnlyAbortedTransaction",
+                        4L,
+                        0,
+                        2,
+                        (RecordProducer) topic -> produceInTransaction(topic, 
0, 3, false)));
+    }
+
+    @ParameterizedTest(name = "{0}")
+    @MethodSource("offsetConvergenceScenarios")
+    void testCommittedOffsetConvergesWithLogEndOffset(
+            String scenario,
+            long expectedLastStableOffset,
+            int expectedReadableRecords,
+            int maxCheckpoints,
+            RecordProducer recordProducer)
+            throws Throwable {
+        final String topic = TOPIC + scenario;
+        final String groupId = scenario + "OffsetCommitGroup";
+        final TopicPartition tp = new TopicPartition(topic, 0);
+        KafkaSourceTestEnv.createTestTopic(topic, 1, 1);
+
+        recordProducer.produce(topic);
+
+        // assert state of Kafka
+        final long lastStableOffset = awaitLastStableOffset(tp, 
expectedLastStableOffset);
+        assertThat(getReadCommittedRecordsCount(tp))
+                .as("Number of readable records")
+                .isEqualTo(expectedReadableRecords);
+
+        // run KafkaSourceReader and check the last offset committed
+        final long committedOffset =
+                readAndCommitOffset(
+                        tp, groupId, expectedReadableRecords, maxCheckpoints, 
lastStableOffset);
+        assertThat(committedOffset)
+                .as("The committed offset should converge with the log end 
offset")
+                .isEqualTo(lastStableOffset);
+    }
+
+    @Test
+    void testCommittedOffsetDoesNotOverrunRecordsInATransaction() throws 
Throwable {
+        final String topic = TOPIC + "PartialTransaction";
+        final String groupId = "PartialTransactionOffsetCommitGroup";
+        final TopicPartition tp = new TopicPartition(topic, 0);
+        KafkaSourceTestEnv.createTestTopic(topic, 1, 1);
+
+        // transaction containing three records
+        // offset 0 = record
+        // offset 1 = record
+        // offset 2 = record
+        // offset 3 = commit marker
+        // offset 4 << log end offset
+        produceInTransaction(topic, 0, 3, true);
+
+        // assert state of Kafka
+        awaitLastStableOffset(tp, 4L);
+        final int numReadableOffsets = getReadCommittedRecordsCount(tp);
+        assertThat(numReadableOffsets).as("Three records should be 
readable").isEqualTo(3);
+
+        // run KafkaSourceReader and check the last offset committed
+        final int recordsToRead = 2;
+        final long committedOffset =
+                readAndCommitOffset(tp, groupId, recordsToRead, 1, 
Long.MIN_VALUE);
+        assertThat(committedOffset)
+                .as("The committed offset should be the next offset after the 
two read records")
+                .isEqualTo(2);
+    }
+
+    /**
+     * Simulates a low-traffic bursty transactional topic, where a burst of 
records are followed by
+     * a commit marker. The result of this could be for a commit marker to be 
returned with no
+     * records in a poll. The intent of this test is to verify that the offset 
converges in such a
+     * scenario, and does not get stuck with an incorrect lag until more 
records arrive on the topic
+     * to allow an update.
+     */
+    @Test
+    void 
testCommittedOffsetConvergesOnIdlePartitionAfterTrailingCommitMarker() throws 
Throwable {
+        final String topic = TOPIC + "IdleAfterTrailingMarker";
+        final String groupId = "IdleAfterTrailingMarkerOffsetCommitGroup";
+        final TopicPartition tp = new TopicPartition(topic, 0);
+        KafkaSourceTestEnv.createTestTopic(topic, 1, 1);
+
+        // transaction containing two records
+        // offset 0 = record
+        // offset 1 = record
+        // offset 2 = commit marker
+        // offset 3 << log end offset
+        produceInTransaction(topic, 0, 2, true);
+
+        // assert state of Kafka
+        final long lastStableOffset = awaitLastStableOffset(tp, 3L);
+        assertThat(getReadCommittedRecordsCount(tp)).as("Number of readable 
records").isEqualTo(2);
+
+        // one record per poll, so that the commit marker is left to a poll of 
its own
+        final Properties oneRecordPerPoll = new Properties();
+        oneRecordPerPoll.setProperty(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 
"1");
+
+        // the poll that skips the commit marker is not the poll that delivers 
the last
+        // record, so allow for the first checkpoint being completed in 
between the two
+        final int maxCheckpoints = 2;
+
+        // run KafkaSourceReader and check the last offset committed
+        final long committedOffset =
+                readAndCommitOffset(
+                        tp, groupId, 2, maxCheckpoints, lastStableOffset, 
oneRecordPerPoll);
+        assertThat(committedOffset)
+                .as("The committed offset should converge with the log end 
offset")
+                .isEqualTo(lastStableOffset);
+    }
+
+    @Test
+    void testCommittedOffsetForBoundedSplitDoesNotGoBackwards() throws 
Throwable {
+        final String topic = TOPIC + "BoundedTrailingMarkers";
+        final String groupId = "BoundedTrailingMarkersOffsetCommitGroup";
+        final TopicPartition tp = new TopicPartition(topic, 0);
+        // second partition to keep the fetcher polling after the bounded 
split finishes
+        final TopicPartition idleTp = new TopicPartition(topic, 1);
+        KafkaSourceTestEnv.createTestTopic(topic, 2, 1);
+
+        // committed transaction containing two records
+        // offset 0 = record
+        // offset 1 = record
+        // offset 2 = commit marker
+        // offset 3 << log end offset
+        produceInTransaction(topic, 0, 2, true);
+        awaitLastStableOffset(tp, 3L);
+
+        final Properties props = new Properties();
+        props.setProperty(ConsumerConfig.GROUP_ID_CONFIG, groupId);
+        props.setProperty(ConsumerConfig.ISOLATION_LEVEL_CONFIG, 
"read_committed");
+
+        final AtomicBoolean splitFinished = new AtomicBoolean(false);
+        try (KafkaSourceReader<Integer> reader =
+                (KafkaSourceReader<Integer>)
+                        createReader(
+                                Boundedness.CONTINUOUS_UNBOUNDED,
+                                new TestingReaderContext(),
+                                (ignore) -> splitFinished.set(true),
+                                props,
+                                null)) {
+            reader.addSplits(
+                    Arrays.asList(
+                            new KafkaPartitionSplit(tp, 0L, 10L),
+                            new KafkaPartitionSplit(
+                                    idleTp, 0L, 
KafkaPartitionSplit.NO_STOPPING_OFFSET)));
+
+            final TestingReaderOutput<Integer> output = new 
TestingReaderOutput<>();
+            pollUntil(
+                    reader,
+                    output,
+                    () -> output.getEmittedRecords().size() == 2,
+                    "The reader did not emit all records before timeout.");
+
+            completeCheckpoint(reader, output, 1L);
+            final long offsetWhileRunning = getCommittedOffset(tp, groupId);
+
+            // aborted transaction that takes the partition up to the split's 
stopping offset
+            // offsets 3 - 8 = records
+            //     offset  9 = abort marker
+            //     offset 10 = log end offset
+            produceInTransaction(topic, 0, 6, false);
+            awaitLastStableOffset(tp, 10L);
+            pollUntil(
+                    reader, output, splitFinished::get, "The split did not 
finish before timeout.");
+
+            completeCheckpoint(reader, output, 2L);
+            final long offsetAfterFinishing = getCommittedOffset(tp, groupId);
+
+            assertThat(offsetAfterFinishing)

Review Comment:
   Yeah, that's an improvement - I've done that now



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