Skip to content
Merged
Show file tree
Hide file tree
Changes from 9 commits
Commits
Show all changes
26 commits
Select commit Hold shift + click to select a range
e894d84
Redo verification path
jolshan Nov 16, 2023
2a52318
Fix build issues
jolshan Nov 17, 2023
865078b
Fix tests
jolshan Nov 17, 2023
a20f238
Merge branch 'trunk' of github.com:apache/kafka into kafka-15784
jolshan Nov 17, 2023
9bac9a9
Fix test failures
jolshan Nov 17, 2023
b4a920b
Rewrite GroupCoordinator and GroupMetadataManager to handle checks fo…
jolshan Nov 29, 2023
05c0d28
Merge branch 'trunk' of github.com:apache/kafka into kafka-15784
jolshan Nov 29, 2023
7bc6d06
Update comments and method names
jolshan Nov 30, 2023
a471f04
Fix style issues and passing verification guards
jolshan Dec 4, 2023
7c01682
Clean up GroupCoordinator, GroupMetadataManger, and ReplicaManager code
jolshan Dec 6, 2023
965ec39
remove package private scoping that is not needed
jolshan Dec 6, 2023
8ecbe4d
Remove produce path refactor
jolshan Dec 7, 2023
99e8577
spacing cleanups
jolshan Dec 7, 2023
a83b136
space
jolshan Dec 7, 2023
56a44ba
remove simple fix
jolshan Dec 7, 2023
e6a14ca
cleanups
jolshan Dec 7, 2023
de5dfad
Update comments simplify locking
jolshan Dec 8, 2023
66fff82
applying cherrypick
jolshan Dec 9, 2023
046425c
Fix tests
jolshan Dec 9, 2023
2c1f6a9
add error conversion
jolshan Dec 9, 2023
3f6c044
add test
jolshan Dec 9, 2023
3a3160d
Fix error code
jolshan Dec 11, 2023
4f757e2
Fix other error
jolshan Dec 11, 2023
ac80ae0
More error fixes and other cleanups
jolshan Dec 12, 2023
13dd654
fix silly error mix up
jolshan Dec 12, 2023
7cad567
rename method
jolshan Dec 12, 2023
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
93 changes: 63 additions & 30 deletions core/src/main/scala/kafka/coordinator/group/GroupCoordinator.scala
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ package kafka.coordinator.group
import java.util.{OptionalInt, Properties}
import java.util.concurrent.atomic.AtomicBoolean
import kafka.common.OffsetAndMetadata
import kafka.server.ReplicaManager.TransactionVerificationEntries
import kafka.server._
import kafka.utils.Logging
import org.apache.kafka.common.{TopicIdPartition, TopicPartition}
Expand Down Expand Up @@ -909,8 +910,68 @@ private[group] class GroupCoordinator(
val group = groupManager.getGroup(groupId).getOrElse {
groupManager.addGroup(new GroupMetadata(groupId, Empty, time))
}
doTxnCommitOffsets(group, transactionalId, memberId, groupInstanceId, generationId, producerId, producerEpoch,
offsetMetadata, requestLocal, responseCallback)

val filteredOffsetMetadata = offsetMetadata.filter { case (_, offsetAndMetadata) =>
groupManager.validateOffsetMetadataLength(offsetAndMetadata.metadata)
}
if (filteredOffsetMetadata.isEmpty) {
// compute the final error codes for the commit response
val commitStatus = offsetMetadata.map { case (k, _) => k -> Errors.OFFSET_METADATA_TOO_LARGE }
responseCallback(commitStatus)
return
}

val magicOpt = groupManager.getMagic(partitionFor(group.groupId))
if (magicOpt.isEmpty) {
val commitStatus = offsetMetadata.map { case (topicIdPartition, _) =>
(topicIdPartition, Errors.NOT_COORDINATOR)
}
responseCallback(commitStatus)
return
}

val records = groupManager.generateOffsetRecords(magicOpt.get, true, group.groupId, filteredOffsetMetadata, producerId, producerEpoch)
val transactionVerificationEntries = new TransactionVerificationEntries

def postVerificationCallback(newRequestLocal: RequestLocal)
(errorResults: Map[TopicPartition, LogAppendResult]): Unit = {
group.inLock {
val validationErrorOpt = validateOffsetCommit(
group,
generationId,
memberId,
groupInstanceId,
isTransactional = true
)

val verifiedOffsets = offsetMetadata.filter {
case (tp, _) =>
!errorResults.contains(tp.topicPartition)
}

val verifiedRecords = records.filter {
case (tp, _) =>
!errorResults.contains(tp)
}

if (validationErrorOpt.isDefined) {
responseCallback(offsetMetadata.map { case (k, _) => k -> validationErrorOpt.get })
} else if (verifiedOffsets.isEmpty) {
responseCallback(offsetMetadata.map { case (k, _) => k -> errorResults(k.topicPartition).error })
} else {
val putCacheCallback = groupManager.createPutCacheCallback(true, group, memberId, offsetMetadata, verifiedOffsets, responseCallback, producerId, verifiedRecords, errorResults)
groupManager.storeOffsetsAfterVerification(group, verifiedOffsets, records, putCacheCallback, producerId, transactionVerificationEntries, errorResults, newRequestLocal)
Comment thread
jolshan marked this conversation as resolved.
Outdated
}
}
}

groupManager.replicaManager.appendRecordsWithTransactionVerification(
entriesPerPartition = records,
transactionVerificationEntries = transactionVerificationEntries,
transactionalId = transactionalId,
requestLocal = requestLocal,
postVerificationCallback = postVerificationCallback
)
}
}

Expand Down Expand Up @@ -951,34 +1012,6 @@ private[group] class GroupCoordinator(
groupManager.scheduleHandleTxnCompletion(producerId, offsetsPartitions.map(_.partition).toSet, isCommit)
}

private def doTxnCommitOffsets(group: GroupMetadata,
transactionalId: String,
memberId: String,
groupInstanceId: Option[String],
generationId: Int,
producerId: Long,
producerEpoch: Short,
offsetMetadata: immutable.Map[TopicIdPartition, OffsetAndMetadata],
requestLocal: RequestLocal,
responseCallback: immutable.Map[TopicIdPartition, Errors] => Unit): Unit = {
group.inLock {
val validationErrorOpt = validateOffsetCommit(
group,
generationId,
memberId,
groupInstanceId,
isTransactional = true
)

if (validationErrorOpt.isDefined) {
responseCallback(offsetMetadata.map { case (k, _) => k -> validationErrorOpt.get })
} else {
groupManager.storeOffsets(group, memberId, offsetMetadata, responseCallback, transactionalId, producerId,
producerEpoch, requestLocal)
}
}
}

private def validateOffsetCommit(
group: GroupMetadata,
generationId: Int,
Expand Down
Loading