From 5dffd47f3da27625db313927e11ce9ae5657c486 Mon Sep 17 00:00:00 2001 From: Andrew Schofield Date: Thu, 28 Mar 2024 18:35:52 +0000 Subject: [PATCH 1/2] AKCORE-81: Experiment with read-committed share groups --- .../main/resources/common/message/ShareFetchResponse.json | 7 +++++++ .../common/message/ShareGroupHeartbeatResponse.json | 2 ++ 2 files changed, 9 insertions(+) diff --git a/clients/src/main/resources/common/message/ShareFetchResponse.json b/clients/src/main/resources/common/message/ShareFetchResponse.json index c1782ec5b23fd..e4e13b660ec98 100644 --- a/clients/src/main/resources/common/message/ShareFetchResponse.json +++ b/clients/src/main/resources/common/message/ShareFetchResponse.json @@ -52,6 +52,13 @@ "The last offset of this batch of acquired records." }, { "name": "DeliveryCount", "type": "int16", "versions": "0+", "about": "The delivery count of this batch of acquired records." } + ]}, + { "name": "AbortedTransactions", "type": "[]AbortedTransaction", "versions": "0+", "nullableVersions": "0+", "ignorable": true, + "about": "The aborted transactions.", "fields": [ + { "name": "ProducerId", "type": "int64", "versions": "0+", "entityType": "producerId", + "about": "The producer id associated with the aborted transaction." }, + { "name": "FirstOffset", "type": "int64", "versions": "0+", + "about": "The first offset in the aborted transaction." } ]} ]} ]}, diff --git a/clients/src/main/resources/common/message/ShareGroupHeartbeatResponse.json b/clients/src/main/resources/common/message/ShareGroupHeartbeatResponse.json index 8a8cbc8e55e56..6c22fa42f1c84 100644 --- a/clients/src/main/resources/common/message/ShareGroupHeartbeatResponse.json +++ b/clients/src/main/resources/common/message/ShareGroupHeartbeatResponse.json @@ -40,6 +40,8 @@ "about": "The member epoch." }, { "name": "HeartbeatIntervalMs", "type": "int32", "versions": "0+", "about": "The heartbeat interval in milliseconds." }, + { "name": "IsolationLevel", "type": "int8", "versions": "0+", "default": "0", "ignorable": true, + "about": "This setting controls the visibility of transactional records. Using READ_UNCOMMITTED (isolation_level = 0) makes all records visible. With READ_COMMITTED (isolation_level = 1), non-transactional and COMMITTED transactional records are visible. To be more concrete, READ_COMMITTED returns all data from offsets smaller than the current LSO (last stable offset), and enables the inclusion of the list of aborted transactions in the result, which allows consumers to discard ABORTED transactional records" }, { "name": "Assignment", "type": "Assignment", "versions": "0+", "nullableVersions": "0+", "default": "null", "about": "null if not provided; the assignment otherwise.", "fields": [ { "name": "Error", "type": "int8", "versions": "0+", From 9360f4c172b3b67b4a16b9ca2dc860a86906ad60 Mon Sep 17 00:00:00 2001 From: Andrew Schofield Date: Wed, 3 Apr 2024 10:36:00 +0100 Subject: [PATCH 2/2] AKCORE-81: Add aborted transactions to ShareFetch response --- .../internals/ShareCompletedFetch.java | 64 +++++++++++++++++++ .../internals/ShareFetchRequestManager.java | 10 ++- .../internals/ShareMemberStateListener.java | 43 +++++++++++++ .../internals/ShareMembershipManager.java | 30 +++++---- .../internals/ShareCompletedFetchTest.java | 2 + .../internals/ShareFetchBufferTest.java | 2 + .../internals/ShareFetchCollectorTest.java | 2 + .../ShareFetchRequestManagerTest.java | 2 +- .../internals/ShareMembershipManagerTest.java | 19 +++--- .../kafka/server/SharePartitionManager.java | 13 ++++ .../main/scala/kafka/server/KafkaApis.scala | 16 ++++- .../test/api/PlaintextShareConsumerTest.java | 47 ++++++++++++++ .../group/GroupMetadataManager.java | 5 ++ 13 files changed, 231 insertions(+), 24 deletions(-) create mode 100644 clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareMemberStateListener.java diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareCompletedFetch.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareCompletedFetch.java index 73b4deaf7a085..cf6cab78a43da 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareCompletedFetch.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareCompletedFetch.java @@ -18,6 +18,7 @@ import org.apache.kafka.clients.consumer.AcknowledgeType; import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.common.IsolationLevel; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.TopicIdPartition; import org.apache.kafka.common.errors.CorruptRecordException; @@ -26,6 +27,7 @@ import org.apache.kafka.common.header.Headers; import org.apache.kafka.common.header.internals.RecordHeaders; import org.apache.kafka.common.message.ShareFetchResponseData; +import org.apache.kafka.common.record.ControlRecordType; import org.apache.kafka.common.record.Record; import org.apache.kafka.common.record.RecordBatch; import org.apache.kafka.common.record.TimestampType; @@ -39,11 +41,15 @@ import java.io.Closeable; import java.nio.ByteBuffer; +import java.util.Comparator; +import java.util.HashSet; import java.util.Iterator; import java.util.LinkedList; import java.util.List; import java.util.ListIterator; import java.util.Optional; +import java.util.PriorityQueue; +import java.util.Set; /** * {@link ShareCompletedFetch} represents a {@link RecordBatch batch} of {@link Record records} @@ -56,11 +62,14 @@ public class ShareCompletedFetch { final TopicIdPartition partition; final ShareFetchResponseData.PartitionData partitionData; + final IsolationLevel isolationLevel; final short requestVersion; private final Logger log; private final BufferSupplier decompressionBufferSupplier; private final Iterator batches; + private final Set abortedProducerIds; + private final PriorityQueue abortedTransactions; private RecordBatch currentBatch; private Record lastRecord; private CloseableIterator records; @@ -75,14 +84,18 @@ public class ShareCompletedFetch { final BufferSupplier decompressionBufferSupplier, final TopicIdPartition partition, final ShareFetchResponseData.PartitionData partitionData, + final IsolationLevel isolationLevel, final short requestVersion) { this.log = logContext.logger(org.apache.kafka.clients.consumer.internals.ShareCompletedFetch.class); this.decompressionBufferSupplier = decompressionBufferSupplier; this.partition = partition; this.partitionData = partitionData; + this.isolationLevel = isolationLevel; this.requestVersion = requestVersion; this.batches = ShareFetchResponse.recordsOrFail(partitionData).batches().iterator(); this.acquiredRecordList = buildAcquiredRecordList(partitionData.acquiredRecords()); + this.abortedProducerIds = new HashSet<>(); + this.abortedTransactions = abortedTransactions(partitionData); } private List buildAcquiredRecordList(List partitionAcquiredRecords) { @@ -302,6 +315,20 @@ private Record nextFetchedRecord(final boolean checkCrcs) { currentBatch = batches.next(); maybeEnsureValid(currentBatch, checkCrcs); + if (isolationLevel == IsolationLevel.READ_COMMITTED && currentBatch.hasProducerId()) { + consumeAbortedTransactionsUpTo(currentBatch.lastOffset()); + + long producerId = currentBatch.producerId(); + if (containsAbortMarker(currentBatch)) { + abortedProducerIds.remove(producerId); + } else if (isBatchAborted(currentBatch)) { + log.debug("Skipping aborted record batch from partition {} with producerId {} and " + + "offsets {} to {}", + partition, producerId, currentBatch.baseOffset(), currentBatch.lastOffset()); + continue; + } + } + records = currentBatch.streamingIterator(decompressionBufferSupplier); } else { Record record = records.next(); @@ -348,6 +375,43 @@ private void maybeCloseRecordStream() { } } + private void consumeAbortedTransactionsUpTo(long offset) { + if (abortedTransactions == null) + return; + + while (!abortedTransactions.isEmpty() && abortedTransactions.peek().firstOffset() <= offset) { + ShareFetchResponseData.AbortedTransaction abortedTransaction = abortedTransactions.poll(); + abortedProducerIds.add(abortedTransaction.producerId()); + } + } + + private boolean isBatchAborted(RecordBatch batch) { + return batch.isTransactional() && abortedProducerIds.contains(batch.producerId()); + } + + private PriorityQueue abortedTransactions(ShareFetchResponseData.PartitionData partition) { + if (partition.abortedTransactions() == null || partition.abortedTransactions().isEmpty()) + return null; + + PriorityQueue abortedTransactions = new PriorityQueue<>( + partition.abortedTransactions().size(), Comparator.comparingLong(ShareFetchResponseData.AbortedTransaction::firstOffset) + ); + abortedTransactions.addAll(partition.abortedTransactions()); + return abortedTransactions; + } + + private boolean containsAbortMarker(RecordBatch batch) { + if (!batch.isControlBatch()) + return false; + + Iterator batchIterator = batch.iterator(); + if (!batchIterator.hasNext()) + return false; + + Record firstRecord = batchIterator.next(); + return ControlRecordType.ABORT == ControlRecordType.parse(firstRecord.key()); + } + private static class OffsetAndDeliveryCount { final long offset; final short deliveryCount; diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareFetchRequestManager.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareFetchRequestManager.java index 137bbddef326c..d82011623affe 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareFetchRequestManager.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareFetchRequestManager.java @@ -20,6 +20,7 @@ import org.apache.kafka.clients.consumer.internals.NetworkClientDelegate.PollResult; import org.apache.kafka.clients.consumer.internals.NetworkClientDelegate.UnsentRequest; import org.apache.kafka.common.Cluster; +import org.apache.kafka.common.IsolationLevel; import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicIdPartition; import org.apache.kafka.common.TopicPartition; @@ -51,7 +52,7 @@ * represent the {@link SubscriptionState#fetchablePartitions(Predicate)} based on the share group * consumer's assignment. */ -public class ShareFetchRequestManager implements RequestManager, MemberStateListener { +public class ShareFetchRequestManager implements RequestManager, ShareMemberStateListener { private final Logger log; private final LogContext logContext; @@ -65,6 +66,7 @@ public class ShareFetchRequestManager implements RequestManager, MemberStateList private final FetchMetricsManager metricsManager; private final IdempotentCloser idempotentCloser = new IdempotentCloser(); private Uuid memberId; + private IsolationLevel isolationLevel = IsolationLevel.READ_UNCOMMITTED; ShareFetchRequestManager(final LogContext logContext, final String groupId, @@ -215,6 +217,7 @@ private void handleShareFetchSuccess(Node fetchTarget, BufferSupplier.create(), partition, partitionData, + isolationLevel, requestVersion); shareFetchBuffer.add(completedFetch); shareFetchBuffer.handleAcknowledgementResponses(partition, Errors.forCode(partitionData.acknowledgeErrorCode())); @@ -333,13 +336,16 @@ public void close() { } @Override - public void onMemberEpochUpdated(Optional memberEpochOpt, Optional memberIdOpt) { + public void onMemberEpochUpdated(Optional memberEpochOpt, Optional memberIdOpt, Optional isolationLevelOpt) { // Only set the memberID once for now - will handle changes in AKCORE-57 if (memberId == null) { if (memberIdOpt.isPresent()) { memberId = Uuid.fromString(memberIdOpt.get()); } } + if (isolationLevelOpt.isPresent()) { + isolationLevel = isolationLevelOpt.get(); + } } @FunctionalInterface diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareMemberStateListener.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareMemberStateListener.java new file mode 100644 index 0000000000000..a3b2fd3e9ec53 --- /dev/null +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareMemberStateListener.java @@ -0,0 +1,43 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.kafka.clients.consumer.internals; + +import org.apache.kafka.common.IsolationLevel; + +import java.util.Optional; + +/** + * Listener for getting notified of member ID and epoch changes. + */ +public interface ShareMemberStateListener { + + /** + * Called whenever member ID or epoch change with new values received from the broker or + * cleared if the member is not part of the group anymore (when it gets fenced, leaves the + * group or fails). + * + * @param memberEpoch New member epoch received from the broker. Empty if the member is + * not part of the group anymore. + * @param memberId Current member ID. Empty if the member is not part of the group. + * @param isolationLevel Isolation level of the group. Empty if the member is not part + * of the group. + */ + void onMemberEpochUpdated(Optional memberEpoch, + Optional memberId, + Optional isolationLevel); +} diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareMembershipManager.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareMembershipManager.java index e965491c21f4e..9e9a276f63d30 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareMembershipManager.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/ShareMembershipManager.java @@ -24,6 +24,7 @@ import org.apache.kafka.clients.consumer.internals.Utils.TopicPartitionComparator; import org.apache.kafka.clients.consumer.internals.events.BackgroundEventHandler; import org.apache.kafka.clients.consumer.internals.events.ErrorBackgroundEvent; +import org.apache.kafka.common.IsolationLevel; import org.apache.kafka.common.TopicIdPartition; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.Uuid; @@ -119,6 +120,12 @@ public class ShareMembershipManager implements RequestManager { */ private String memberId = ""; + /** + * Isolation level of the share group, received in a heartbeat response when joining the + * group specified in {@link #groupId} + */ + private IsolationLevel isolationLevel = IsolationLevel.READ_UNCOMMITTED; + /** * Current epoch of the member. It will be set to 0 by the member, and provided to the server * on the heartbeat request, to join the group. It will be then maintained by the server, @@ -197,7 +204,7 @@ public class ShareMembershipManager implements RequestManager { * values received from the broker, or values cleared due to member leaving the group, getting * fenced or failing). */ - private final List stateUpdatesListeners; + private final List stateUpdatesListeners; /** * Optional client telemetry reporter which sends client telemetry data to the broker. This @@ -306,7 +313,7 @@ public void onHeartbeatResponseReceived(ShareGroupHeartbeatResponseData response } this.memberId = response.memberId(); - updateMemberEpoch(response.memberEpoch()); + updateMemberEpoch(response.memberEpoch(), IsolationLevel.forId(response.isolationLevel())); ShareGroupHeartbeatResponseData.Assignment assignment = response.assignment(); @@ -412,7 +419,7 @@ public void transitionToFatal() { MemberState previousState = state; transitionTo(MemberState.FATAL); log.error("Member {} with epoch {} transitioned to {} state", memberId, memberEpoch, MemberState.FATAL); - notifyEpochChange(Optional.empty(), Optional.empty()); + notifyEpochChange(Optional.empty(), Optional.empty(), Optional.empty()); if (previousState == MemberState.UNSUBSCRIBED) { log.debug("Member {} with epoch {} got fatal error from the broker but it already " + @@ -546,7 +553,7 @@ void transitionToSendingLeaveGroup() { memberId); return; } - updateMemberEpoch(ShareGroupHeartbeatRequest.LEAVE_GROUP_MEMBER_EPOCH); + updateMemberEpoch(ShareGroupHeartbeatRequest.LEAVE_GROUP_MEMBER_EPOCH, IsolationLevel.READ_UNCOMMITTED); currentAssignment = new HashMap<>(); transitionTo(MemberState.LEAVING); } @@ -556,8 +563,8 @@ void transitionToSendingLeaveGroup() { * This also includes the latest member ID in the notification. If the member fails or leaves * the group, this will be invoked with empty epoch and member ID. */ - private void notifyEpochChange(Optional epoch, Optional memberId) { - stateUpdatesListeners.forEach(stateListener -> stateListener.onMemberEpochUpdated(epoch, memberId)); + private void notifyEpochChange(Optional epoch, Optional memberId, Optional isolationLevel) { + stateUpdatesListeners.forEach(stateListener -> stateListener.onMemberEpochUpdated(epoch, memberId, isolationLevel)); } /** @@ -956,19 +963,20 @@ private void clearPendingAssignmentsAndLocalNamesCache() { } private void resetEpoch() { - updateMemberEpoch(ShareGroupHeartbeatRequest.JOIN_GROUP_MEMBER_EPOCH); + updateMemberEpoch(ShareGroupHeartbeatRequest.JOIN_GROUP_MEMBER_EPOCH, IsolationLevel.READ_UNCOMMITTED); } - private void updateMemberEpoch(int newEpoch) { + private void updateMemberEpoch(int newEpoch, IsolationLevel newIsolationLevel) { boolean newEpochReceived = this.memberEpoch != newEpoch; this.memberEpoch = newEpoch; + isolationLevel = newIsolationLevel; // Simply notify based on epoch change only, given that the member will never receive a // new member ID without an epoch (member ID is only assigned when it joins the group). if (newEpochReceived) { if (memberEpoch > 0) { - notifyEpochChange(Optional.of(memberEpoch), Optional.ofNullable(memberId)); + notifyEpochChange(Optional.of(memberEpoch), Optional.ofNullable(memberId), Optional.of(isolationLevel)); } else { - notifyEpochChange(Optional.empty(), Optional.empty()); + notifyEpochChange(Optional.empty(), Optional.empty(), Optional.empty()); } } } @@ -1038,7 +1046,7 @@ boolean reconciliationInProgress() { * * @param listener Listener to invoke. */ - public void registerStateListener(MemberStateListener listener) { + public void registerStateListener(ShareMemberStateListener listener) { if (listener == null) { throw new IllegalArgumentException("State updates listener cannot be null"); } diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareCompletedFetchTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareCompletedFetchTest.java index dd87da2ff6bc3..5afb185a49d44 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareCompletedFetchTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareCompletedFetchTest.java @@ -18,6 +18,7 @@ import org.apache.kafka.clients.consumer.AcknowledgeType; import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.common.IsolationLevel; import org.apache.kafka.common.TopicIdPartition; import org.apache.kafka.common.Uuid; import org.apache.kafka.common.message.ShareFetchResponseData; @@ -263,6 +264,7 @@ private ShareCompletedFetch newShareCompletedFetch(ShareFetchResponseData.Partit BufferSupplier.create(), TIP, partitionData, + IsolationLevel.READ_UNCOMMITTED, ApiKeys.SHARE_FETCH.latestVersion()); } diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchBufferTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchBufferTest.java index 95d1ba57735a5..3a8af25f5630c 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchBufferTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchBufferTest.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.clients.consumer.internals; +import org.apache.kafka.common.IsolationLevel; import org.apache.kafka.common.TopicIdPartition; import org.apache.kafka.common.Uuid; import org.apache.kafka.common.message.ShareFetchResponseData; @@ -152,6 +153,7 @@ private ShareCompletedFetch completedFetch(TopicIdPartition tp) { BufferSupplier.create(), tp, partitionData, + IsolationLevel.READ_UNCOMMITTED, ApiKeys.SHARE_FETCH.latestVersion()); } diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchCollectorTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchCollectorTest.java index df405e7f4adc0..e07e221b62e71 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchCollectorTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchCollectorTest.java @@ -17,6 +17,7 @@ package org.apache.kafka.clients.consumer.internals; import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.common.IsolationLevel; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.TopicIdPartition; import org.apache.kafka.common.Uuid; @@ -335,6 +336,7 @@ private ShareCompletedFetch build() { BufferSupplier.create(), topicAPartition0, partitionData, + IsolationLevel.READ_UNCOMMITTED, ApiKeys.SHARE_FETCH.latestVersion()); } } diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchRequestManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchRequestManagerTest.java index d79b01d803060..17e8f584a73ad 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchRequestManagerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareFetchRequestManagerTest.java @@ -637,7 +637,7 @@ public TestableShareFetchRequestManager(LogContext logContext, ShareFetchCollector fetchCollector) { super(logContext, groupId, metadata, subscriptions, fetchConfig, shareFetchBuffer, metricsManager); this.shareFetchCollector = fetchCollector; - onMemberEpochUpdated(Optional.empty(), Optional.of(Uuid.randomUuid().toString())); + onMemberEpochUpdated(Optional.empty(), Optional.of(Uuid.randomUuid().toString()), Optional.of(IsolationLevel.READ_UNCOMMITTED)); } private ShareFetch collectFetch() { diff --git a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareMembershipManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareMembershipManagerTest.java index 863112b9b21fa..95a16a6ebe23e 100644 --- a/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareMembershipManagerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/consumer/internals/ShareMembershipManagerTest.java @@ -17,6 +17,7 @@ package org.apache.kafka.clients.consumer.internals; import org.apache.kafka.clients.consumer.internals.events.BackgroundEventHandler; +import org.apache.kafka.common.IsolationLevel; import org.apache.kafka.common.TopicIdPartition; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.Uuid; @@ -207,37 +208,37 @@ public void testFencingWhenStateIsStable() { public void testListenersGetNotifiedOnTransitionsToFatal() { ShareMembershipManager membershipManager = createMembershipManagerJoiningGroup(); subscriptionState.subscribe(Collections.singleton("topic1"), Optional.empty()); - MemberStateListener listener = mock(MemberStateListener.class); + ShareMemberStateListener listener = mock(ShareMemberStateListener.class); membershipManager.registerStateListener(listener); mockStableMember(membershipManager); - verify(listener).onMemberEpochUpdated(Optional.of(MEMBER_EPOCH), Optional.of(MEMBER_ID)); + verify(listener).onMemberEpochUpdated(Optional.of(MEMBER_EPOCH), Optional.of(MEMBER_ID), Optional.of(IsolationLevel.READ_UNCOMMITTED)); clearInvocations(listener); // Transition to FAILED before getting member ID/epoch membershipManager.transitionToFatal(); assertEquals(MemberState.FATAL, membershipManager.state()); - verify(listener).onMemberEpochUpdated(Optional.empty(), Optional.empty()); + verify(listener).onMemberEpochUpdated(Optional.empty(), Optional.empty(), Optional.empty()); } @Test public void testListenersGetNotifiedOnTransitionsToLeavingGroup() { ShareMembershipManager membershipManager = createMembershipManagerJoiningGroup(); - MemberStateListener listener = mock(MemberStateListener.class); + ShareMemberStateListener listener = mock(ShareMemberStateListener.class); membershipManager.registerStateListener(listener); mockStableMember(membershipManager); - verify(listener).onMemberEpochUpdated(Optional.of(MEMBER_EPOCH), Optional.of(MEMBER_ID)); + verify(listener).onMemberEpochUpdated(Optional.of(MEMBER_EPOCH), Optional.of(MEMBER_ID), Optional.of(IsolationLevel.READ_UNCOMMITTED)); clearInvocations(listener); mockLeaveGroup(); membershipManager.leaveGroup(); assertEquals(MemberState.LEAVING, membershipManager.state()); - verify(listener).onMemberEpochUpdated(Optional.empty(), Optional.empty()); + verify(listener).onMemberEpochUpdated(Optional.empty(), Optional.empty(), Optional.empty()); } @Test public void testListenersGetNotifiedOfMemberEpochUpdatesOnlyIfItChanges() { ShareMembershipManager membershipManager = createMembershipManagerJoiningGroup(); - MemberStateListener listener = mock(MemberStateListener.class); + ShareMemberStateListener listener = mock(ShareMemberStateListener.class); membershipManager.registerStateListener(listener); int epoch = 5; @@ -246,14 +247,14 @@ public void testListenersGetNotifiedOfMemberEpochUpdatesOnlyIfItChanges() { .setMemberId(MEMBER_ID) .setMemberEpoch(epoch)); - verify(listener).onMemberEpochUpdated(Optional.of(epoch), Optional.of(MEMBER_ID)); + verify(listener).onMemberEpochUpdated(Optional.of(epoch), Optional.of(MEMBER_ID), Optional.of(IsolationLevel.READ_UNCOMMITTED)); clearInvocations(listener); membershipManager.onHeartbeatResponseReceived(new ShareGroupHeartbeatResponseData() .setErrorCode(Errors.NONE.code()) .setMemberId(MEMBER_ID) .setMemberEpoch(epoch)); - verify(listener, never()).onMemberEpochUpdated(any(), any()); + verify(listener, never()).onMemberEpochUpdated(any(), any(), any()); } private void mockStableMember(ShareMembershipManager membershipManager) { diff --git a/core/src/main/java/kafka/server/SharePartitionManager.java b/core/src/main/java/kafka/server/SharePartitionManager.java index 076805fc8ee5d..2c371371e8ed7 100644 --- a/core/src/main/java/kafka/server/SharePartitionManager.java +++ b/core/src/main/java/kafka/server/SharePartitionManager.java @@ -19,6 +19,7 @@ import org.apache.kafka.common.TopicIdPartition; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.Uuid; +import org.apache.kafka.common.message.FetchResponseData; import org.apache.kafka.common.message.ShareAcknowledgeResponseData; import org.apache.kafka.common.message.ShareFetchResponseData; import org.apache.kafka.common.protocol.Errors; @@ -202,6 +203,18 @@ private Map processFetch .setErrorCode(fetchPartitionData.error.code()) .setAcquiredRecords(acquiredRecords) .setAcknowledgeErrorCode(Errors.NONE.code()); + + if (fetchPartitionData.abortedTransactions.isPresent()) { + List fetchedTransactions = fetchPartitionData.abortedTransactions.get(); + List abortedTransactions = new ArrayList<>(fetchedTransactions.size()); + fetchedTransactions.forEach(fetchedTransaction -> { + ShareFetchResponseData.AbortedTransaction abortedTransaction = new ShareFetchResponseData.AbortedTransaction(); + abortedTransaction.setProducerId(fetchedTransaction.producerId()); + abortedTransaction.setFirstOffset(fetchedTransaction.firstOffset()); + abortedTransactions.add(abortedTransaction); + }); + partitionData.setAbortedTransactions(abortedTransactions); + } } result.put(topicIdPartition, partitionData); }); diff --git a/core/src/main/scala/kafka/server/KafkaApis.scala b/core/src/main/scala/kafka/server/KafkaApis.scala index 2f3c9d9406d07..3c658a00191bb 100644 --- a/core/src/main/scala/kafka/server/KafkaApis.scala +++ b/core/src/main/scala/kafka/server/KafkaApis.scala @@ -1314,6 +1314,7 @@ class KafkaApis(val requestChannel: RequestChannel, .setErrorCode(Errors.forCode(partitionData.errorCode).code) .setRecords(unconvertedRecords) .setAcquiredRecords(partitionData.acquiredRecords) + .setAbortedTransactions(partitionData.abortedTransactions()) .setCurrentLeader(partitionData.currentLeader()) } @@ -1424,7 +1425,7 @@ class KafkaApis(val requestChannel: RequestChannel, request.context.listenerName.value)) // Dummy values for replicaId, replicaEpoch, isolationLevel and clientMetadata as they are not used in ShareFetchRequest - val params = new FetchParams( + var params = new FetchParams( versionId, FetchRequest.FUTURE_LOCAL_REPLICA_ID, -1, @@ -1435,6 +1436,19 @@ class KafkaApis(val requestChannel: RequestChannel, clientMetadata ) + if (groupId.equals("read-committed")) { + params = new FetchParams( + versionId, + FetchRequest.ORDINARY_CONSUMER_ID, + -1, + shareFetchRequest.maxWait, + fetchMinBytes, + fetchMaxBytes, + FetchIsolation.TXN_COMMITTED, + clientMetadata + ) + } + // TODO: We can simplify the code by removing the need for the additional max bytes map. Rather // change "interesting" a tuple to contain both the partition data and the max bytes. val partitionMaxBytes = new util.LinkedHashMap[TopicIdPartition, Integer] diff --git a/core/src/test/java/kafka/test/api/PlaintextShareConsumerTest.java b/core/src/test/java/kafka/test/api/PlaintextShareConsumerTest.java index 91b3a8bd6401b..d20f15586377a 100644 --- a/core/src/test/java/kafka/test/api/PlaintextShareConsumerTest.java +++ b/core/src/test/java/kafka/test/api/PlaintextShareConsumerTest.java @@ -301,6 +301,53 @@ public void testControlRecordsSkipped(String quorum) throws Exception { shareConsumer.close(); } + @ParameterizedTest(name = TEST_WITH_PARAMETERIZED_QUORUM_NAME) + @ValueSource(strings = {"kraft+kip932"}) + public void testReadCommitted(String quorum) throws Exception { + ProducerRecord record = new ProducerRecord<>(tp().topic(), tp().partition(), null, "key".getBytes(), "value".getBytes()); + + Properties transactionProducerProps = new Properties(); + transactionProducerProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "T1"); + KafkaProducer transactionalProducer = createProducer(new ByteArraySerializer(), new ByteArraySerializer(), transactionProducerProps); + transactionalProducer.initTransactions(); + transactionalProducer.beginTransaction(); + RecordMetadata transactional1 = transactionalProducer.send(record).get(); + + KafkaProducer nonTransactionalProducer = createProducer(new ByteArraySerializer(), new ByteArraySerializer(), new Properties()); + RecordMetadata nonTransactional1 = nonTransactionalProducer.send(record).get(); + + transactionalProducer.commitTransaction(); + + // Because this record is in a transaction which is aborted, the consumer should not receive it + transactionalProducer.beginTransaction(); + RecordMetadata transactional2 = transactionalProducer.send(record).get(); + transactionalProducer.abortTransaction(); + + RecordMetadata nonTransactional2 = nonTransactionalProducer.send(record).get(); + + transactionalProducer.close(); + nonTransactionalProducer.close(); + + Properties props = new Properties(); + props.put(ConsumerConfig.GROUP_ID_CONFIG, "read-committed"); + KafkaShareConsumer shareConsumer = createShareConsumer(new ByteArrayDeserializer(), new ByteArrayDeserializer(), + props, CollectionConverters.asScala(Collections.emptyList()).toList()); + shareConsumer.subscribe(Collections.singleton(tp().topic())); + ConsumerRecords records = shareConsumer.poll(Duration.ofMillis(5000000)); + assertEquals(3, records.count()); + assertEquals(transactional1.offset(), records.records(tp()).get(0).offset()); + assertEquals(nonTransactional1.offset(), records.records(tp()).get(1).offset()); + assertEquals(nonTransactional2.offset(), records.records(tp()).get(2).offset()); + + // There will be control records on the topic-partition, so the offsets of the non-control records + // are not 0, 1, 2, 3. Just assert that the offset of the final one is not 3. + assertNotEquals(3, nonTransactional2.offset()); + + records = shareConsumer.poll(Duration.ofMillis(5000)); + assertEquals(0, records.count()); + shareConsumer.close(); + } + @ParameterizedTest(name = TEST_WITH_PARAMETERIZED_QUORUM_NAME) @ValueSource(strings = {"kraft+kip932"}) public void testExplicitAcknowledgeSuccess(String quorum) throws Exception { diff --git a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java index 588ad661259ba..e0950c2b0572b 100644 --- a/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java +++ b/group-coordinator/src/main/java/org/apache/kafka/coordinator/group/GroupMetadataManager.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.coordinator.group; +import org.apache.kafka.common.IsolationLevel; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.Uuid; import org.apache.kafka.common.errors.ApiException; @@ -1575,6 +1576,10 @@ private CoordinatorResult shareGroupHea .setMemberEpoch(updatedMember.memberEpoch()) .setHeartbeatIntervalMs(shareGroupHeartbeatIntervalMs); + if (groupId.equals("read-committed")) { + response.setIsolationLevel(IsolationLevel.READ_COMMITTED.id()); + } + // The assignment is only provided in the following cases: // 1. The member just joined or rejoined to group (epoch equals to zero); // 2. The member's assignment has been updated.