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


##########
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)
+                    .as("The committed offset should not go backwards when the 
split finishes")
+                    .isGreaterThanOrEqualTo(offsetWhileRunning);
+        }
+    }
+
+    /**
+     * Snapshots the reader and waits for the resulting offsets to reach 
Kafka, polling the reader
+     * throughout, for consistency with how the mailbox thread interleaves 
pollNext with checkpoint
+     * callbacks.
+     */
+    private void completeCheckpoint(
+            KafkaSourceReader<Integer> reader, ReaderOutput<Integer> output, 
long checkpointId)
+            throws Exception {
+        try {
+            reader.snapshotState(checkpointId);
+        } catch (Exception e) {
+            throw new RuntimeException(e);
+        }
+        waitUtil(
+                () -> {
+                    try {
+                        reader.notifyCheckpointComplete(checkpointId);
+                        reader.pollNext(output);
+                    } catch (Exception exception) {
+                        throw new RuntimeException(
+                                "Unexpected exception while committing", 
exception);
+                    }
+                    return reader.getOffsetsToCommit().isEmpty();
+                },
+                Duration.ofSeconds(30),
+                Duration.ofMillis(500),
+                "Offset commit did not finish before timeout.");
+    }
+
+    /**
+     * Reads {@code expectedRecords} records from {@code tp} using a 
read_committed reader, and then
+     * completes checkpoints while the partition is idle until either the 
committed offset reaches
+     * {@code targetOffset} or {@code maxCheckpoints} have been completed. 
Returns the last offset
+     * that the reader committed to Kafka.
+     */
+    private long readAndCommitOffset(
+            TopicPartition tp,
+            String groupId,
+            int expectedRecords,
+            int maxCheckpoints,
+            long targetOffset)
+            throws Exception {
+        return readAndCommitOffset(
+                tp, groupId, expectedRecords, maxCheckpoints, targetOffset, 
new Properties());
+    }
+
+    private long readAndCommitOffset(
+            TopicPartition tp,
+            String groupId,
+            int expectedRecords,
+            int maxCheckpoints,
+            long targetOffset,
+            Properties extraConsumerProps)
+            throws Exception {
+        final Properties props = new Properties();
+        props.setProperty(ConsumerConfig.GROUP_ID_CONFIG, groupId);
+        props.setProperty(ConsumerConfig.ISOLATION_LEVEL_CONFIG, 
"read_committed");
+        props.putAll(extraConsumerProps);
+
+        final MetricListener metricListener = new MetricListener();
+
+        try (KafkaSourceReader<Integer> reader =
+                (KafkaSourceReader<Integer>)
+                        createReader(
+                                Boundedness.CONTINUOUS_UNBOUNDED,
+                                new TestingReaderContext(
+                                        new Configuration(),
+                                        InternalSourceReaderMetricGroup.mock(
+                                                
metricListener.getMetricGroup())),
+                                (ignore) -> {},
+                                props,
+                                null)) {
+            // prepare kafka reader
+            reader.addSplits(
+                    Collections.singletonList(
+                            new KafkaPartitionSplit(
+                                    tp, 0L, 
KafkaPartitionSplit.NO_STOPPING_OFFSET)));

Review Comment:
   Yes, that's a good point, that was a gap. A split starting on 
EARLIEST_OFFSET would keep that until it emitted a record, and snapshotState 
was skipping splits with a negative offset. So an aborted-only partition 
wouldn't commit anything. 
   
   I've added this to the convergence scenarios 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