Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 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"Electing as the group coordinator for partition $offsetTopicPartitionId")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Shouldn't that read Elected as the group coordinator ...?

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"Resigning 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"Becoming the txn coordinator for partition $txnTopicPartitionId at epoch $coordinatorEpoch")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do you also plan to change this to Elected as ...?

// 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"Resigning the txn coordinator for partition $txnTopicPartitionId at epoch $coordinatorEpoch")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Same as above.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This still needs to be changed.

coordinatorEpoch match {
case Some(epoch) =>
txnManager.removeTransactionsForTxnTopicPartition(txnTopicPartitionId, epoch)
Expand Down