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]