Savonitar commented on code in PR #293:
URL:
https://github.com/apache/flink-connector-kafka/pull/293#discussion_r4155597664
##########
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/reader/KafkaPartitionSplitReader.java:
##########
@@ -389,6 +400,59 @@ private long getConsumerPosition(
return retryOnWakeup(() -> consumer.position(tp), msg);
}
+ private void trackLastFetchedRecordOffsets(ConsumerRecords<byte[], byte[]>
consumerRecords) {
+ for (TopicPartition tp : consumerRecords.partitions()) {
+ List<ConsumerRecord<byte[], byte[]>> partitionRecords =
consumerRecords.records(tp);
+ if (!partitionRecords.isEmpty()) {
+ lastFetchedOffsets.put(
+ tp, partitionRecords.get(partitionRecords.size() -
1).offset());
+ }
+ }
+ }
+
+ /**
+ * Advances the offsets to commit over the entries that the Kafka consumer
read but never
+ * delivered, such as transaction control markers and records of aborted
transactions.
+ *
+ * <p>{@link KafkaRecordEmitter} derives the offset to commit from the
records it receives, so
+ * the offset stops at the first entry that Kafka does not deliver, for as
long as the partition
+ * is idle. The consumer's own position accounts for those entries, so it
is the offset that
+ * external tooling expects to see.
+ *
+ * <p>This relies on {@link #lastKnownPositions}, populated as a side
effect of the regular
+ * {@link #fetch()} poll loop, rather than querying the consumer for the
position again here.
+ * {@link #notifyCheckpointComplete} runs on this same split fetcher
thread, but at a point
+ * outside that poll loop, so this allows us to avoid a separate blocking
call to the consumer.
+ */
+ private Map<TopicPartition, OffsetAndMetadata> reconcileOffsetsToCommit(
+ KafkaConsumer<byte[], byte[]> consumer,
+ Map<TopicPartition, OffsetAndMetadata> offsetsToCommit) {
+ Map<TopicPartition, OffsetAndMetadata> reconciled = new
HashMap<>(offsetsToCommit);
+ Set<TopicPartition> assignment = consumer.assignment();
+ offsetsToCommit.forEach(
+ (tp, offsetAndMetadata) -> {
+ if (!assignment.contains(tp) ||
stoppingOffsets.containsKey(tp)) {
+ return;
+ }
+ Long lastFetchedOffset = lastFetchedOffsets.get(tp);
+ if (lastFetchedOffset != null
+ && offsetAndMetadata.offset() != lastFetchedOffset
+ 1) {
Review Comment:
Could this move the committed offset backwards when new records arrive
between a checkpoint's snapshot and its commit?
Example (transactional record at 0, commit marker at 1):
1.The readier emits the record at offset 0, consumer moves position (skips
the commit marker at 1) to 2.
2.Checkpoint 1 snapshots offset 1 and reconciliation commits 2.
3.Checkpoint 2 snapshots offset 1 again while the partition is idle.
4.A new record at offset 2 arrives and is consumed before checkpoint 2
completes.
This check now rejects reconciliation because 1 != 2 + 1, so checkpoint 2
commits the snapshotted offset 1.
Kafka's committed offset therefore moves from 2 back to 1 until the next
checkpoint.
The impact is monitoring flick/blip. But
[testCommittedOffsetForBoundedSplitDoesNotGoBackwards](https://github.com/apache/flink-connector-kafka/pull/293/changes#diff-912e48531349dc6ecdd8a8a2c25c5b65d4f3511bdd028a0e1cdb5605350759f3R504)
suggests we want the committed offset to only move **forward**.
Could we preserve the progress already established? Or am I missing
something?
##########
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:
> Thread.sleep(200L);
The class already uses CommonTestUtils.waitUtil for the same purpose, could
we reuse it here?
##########
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:
Could we also run the aborted-only scenario with
KafkaPartitionSplit.EARLIEST_OFFSET, matching the default
OffsetsInitializer.earliest() ?
##########
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:
if the probe hasn't reached the LSO within 30s, the loop just exits and
returns what it has counted so far. The failure then shows up as a misleading:
`Number of readable records expected 3 but was 2`
instead of saying that the probe timed out.
The loop also uses its own deadline rather than waitUtil.
##########
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)));
+
+ // read until partition is idle
+ final TestingReaderOutput<Integer> output = new
TestingReaderOutput<>();
+ pollUntil(
+ reader,
+ output,
+ () -> output.getEmittedRecords().size() == expectedRecords,
+ "The reader did not emit all records before timeout.");
+
+ // complete checkpoints to commit offset
+ long committedOffset = Long.MIN_VALUE;
+ for (int checkpointId = 1; checkpointId <= maxCheckpoints;
checkpointId++) {
+ final long currentCheckpointId = checkpointId;
+ if (checkpointId > 1) {
+ reader.pollNext(output);
+ }
+ reader.snapshotState(currentCheckpointId);
Review Comment:
this block of code looks exactly as
https://github.com/apache/flink-connector-kafka/pull/293/changes#diff-912e48531349dc6ecdd8a8a2c25c5b65d4f3511bdd028a0e1cdb5605350759f3R572
(the method `completeCheckpoint` that u added in this pr).
Could we follow DRY principle here?
##########
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:
did u consider to use the exact result via isEqualTo to harden the test and
document the accepted behavior: bounderd splits do not skip traling marker ?
e.g. maybe
```
assertThat(offsetWhileRunning)
.as("The committed offset of a running bounded split")
.isEqualTo(2L);
```
##########
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 {
Review Comment:
Could you please clarify, why do we need this try-catch?
--
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]