dalelane commented on code in PR #293:
URL:
https://github.com/apache/flink-connector-kafka/pull/293#discussion_r4161781220
##########
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:
This was right - the check was (correctly) refusing to advance past a record
that the checkpoint hasn't seen. But it was falling back to the raw snapshot
offset, discarding the progress already committed for the same offset.
I've added a
`testCommittedOffsetDoesNotGoBackwardsWhenRecordIsFetchedBeforeCommit` test for
this scenario. (The split reader now remembers the offset it last committed for
each assigned partition. It leaves a partition out of a commit rather than
committing a lower offset.)
--
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]