From 76ff0696fda318d637cf397bbe6dacc7e6bec25f Mon Sep 17 00:00:00 2001 From: dengziming Date: Thu, 2 Jun 2022 11:25:30 +0800 Subject: [PATCH] MINOR: Use right enum value for broker registration change --- .../controller/ClusterControlManager.java | 4 +- .../controller/ReplicationControlManager.java | 4 +- .../org/apache/kafka/image/ClusterDelta.java | 4 +- .../BrokerRegistrationFencingChange.java | 4 +- .../controller/ClusterControlManagerTest.java | 6 +-- .../BrokerRegistrationFencingChangeTest.java | 8 ++-- ...trationInControlledShutdownChangeTest.java | 48 +++++++++++++++++++ 7 files changed, 63 insertions(+), 15 deletions(-) create mode 100644 metadata/src/test/java/org/apache/kafka/metadata/BrokerRegistrationInControlledShutdownChangeTest.java diff --git a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java index 859dedda2200f..ad11e19ec2d74 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ClusterControlManager.java @@ -422,7 +422,7 @@ public void replay(FenceBrokerRecord record) { record, record.id(), record.epoch(), - BrokerRegistrationFencingChange.UNFENCE.asBoolean(), + BrokerRegistrationFencingChange.FENCE.asBoolean(), BrokerRegistrationInControlledShutdownChange.NONE.asBoolean() ); } @@ -432,7 +432,7 @@ public void replay(UnfenceBrokerRecord record) { record, record.id(), record.epoch(), - BrokerRegistrationFencingChange.FENCE.asBoolean(), + BrokerRegistrationFencingChange.UNFENCE.asBoolean(), BrokerRegistrationInControlledShutdownChange.NONE.asBoolean() ); } diff --git a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java index 8382cd9c16b7b..fc6ac38c31728 100644 --- a/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java +++ b/metadata/src/main/java/org/apache/kafka/controller/ReplicationControlManager.java @@ -1162,7 +1162,7 @@ void handleBrokerFenced(int brokerId, List records) { if (featureControl.metadataVersion().isBrokerRegistrationChangeRecordSupported()) { records.add(new ApiMessageAndVersion(new BrokerRegistrationChangeRecord(). setBrokerId(brokerId).setBrokerEpoch(brokerRegistration.epoch()). - setFenced(BrokerRegistrationFencingChange.UNFENCE.value()), + setFenced(BrokerRegistrationFencingChange.FENCE.value()), (short) 0)); } else { records.add(new ApiMessageAndVersion(new FenceBrokerRecord(). @@ -1205,7 +1205,7 @@ void handleBrokerUnfenced(int brokerId, long brokerEpoch, List