Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
Original file line number Diff line number Diff line change
Expand Up @@ -908,6 +908,7 @@ class GroupCoordinator(val brokerId: Int,
* @param offsetTopicPartitionId The partition we are now leading
*/
def onElection(offsetTopicPartitionId: Int): Unit = {
info(s"Elected as the group coordinator for partition $offsetTopicPartitionId")
groupManager.scheduleLoadGroupAndOffsets(offsetTopicPartitionId, onGroupLoaded)
}

Expand All @@ -917,6 +918,7 @@ class GroupCoordinator(val brokerId: Int,
* @param offsetTopicPartitionId The partition we are no longer leading
*/
def onResignation(offsetTopicPartitionId: Int): Unit = {
info(s"Resigned as the group coordinator for partition $offsetTopicPartitionId")
groupManager.removeGroupsForPartition(offsetTopicPartitionId, onGroupUnloaded)
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -544,6 +544,7 @@ class GroupMetadataManager(brokerId: Int,
private[group] def loadGroupsAndOffsets(topicPartition: TopicPartition, onGroupLoaded: GroupMetadata => Unit, startTimeMs: java.lang.Long): Unit = {
try {
val schedulerTimeMs = time.milliseconds() - startTimeMs
debug(s"Started loading offsets and group metadata from $topicPartition")
doLoadGroupsAndOffsets(topicPartition, onGroupLoaded)
val endTimeMs = time.milliseconds()
val totalLoadingTimeMs = endTimeMs - startTimeMs
Expand Down Expand Up @@ -759,6 +760,7 @@ class GroupMetadataManager(brokerId: Int,
var numOffsetsRemoved = 0
var numGroupsRemoved = 0

debug(s"Started unloading offsets and group metadata for $topicPartition")
inLock(partitionLock) {
// we need to guard the group removal in cache in the loading partition lock
// to prevent coordinator's check-and-get-group race condition
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -311,6 +311,7 @@ class TransactionCoordinator(brokerId: Int,
* @param coordinatorEpoch The partition coordinator (or leader) epoch from the received LeaderAndIsr request
*/
def onElection(txnTopicPartitionId: Int, coordinatorEpoch: Int): Unit = {
info(s"Elected as the txn coordinator for partition $txnTopicPartitionId at epoch $coordinatorEpoch")
// The operations performed during immigration must be resilient to any previous errors we saw or partial state we
// left off during the unloading phase. Ensure we remove all associated state for this partition before we continue
// loading it.
Expand All @@ -329,6 +330,7 @@ class TransactionCoordinator(brokerId: Int,
* are resigning after receiving a StopReplica request from the controller
*/
def onResignation(txnTopicPartitionId: Int, coordinatorEpoch: Option[Int]): Unit = {
info(s"Resigned as the txn coordinator for partition $txnTopicPartitionId at epoch $coordinatorEpoch")
coordinatorEpoch match {
case Some(epoch) =>
txnManager.removeTransactionsForTxnTopicPartition(txnTopicPartitionId, epoch)
Expand Down