diff --git a/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java b/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java index 021730cfb3..410092dfd4 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java @@ -1006,7 +1006,7 @@ public class BrokerOuterAPI { final RemotingCommand response = this.remotingClient.invokeSync(controllerAddress, request, 3000); assert response != null; if (response.getCode() == SUCCESS) { - return response.decodeCommandCustomHeader(GetMetaDataResponseHeader.class); + return (GetMetaDataResponseHeader) response.decodeCommandCustomHeader(GetMetaDataResponseHeader.class); } throw new MQBrokerException(response.getCode(), response.getRemark()); } @@ -1050,7 +1050,7 @@ public class BrokerOuterAPI { assert response != null; switch (response.getCode()) { case SUCCESS: { - return response.decodeCommandCustomHeader(RegisterBrokerToControllerResponseHeader.class); + return (RegisterBrokerToControllerResponseHeader) response.decodeCommandCustomHeader(RegisterBrokerToControllerResponseHeader.class); } case CONTROLLER_NOT_LEADER: { throw new MQBrokerException(response.getCode(), "Controller leader was changed"); @@ -1070,7 +1070,7 @@ public class BrokerOuterAPI { assert response != null; switch (response.getCode()) { case SUCCESS: { - final GetReplicaInfoResponseHeader header = response.decodeCommandCustomHeader(GetReplicaInfoResponseHeader.class); + final GetReplicaInfoResponseHeader header = (GetReplicaInfoResponseHeader) response.decodeCommandCustomHeader(GetReplicaInfoResponseHeader.class); assert response.getBody() != null; final SyncStateSet stateSet = RemotingSerializable.decode(response.getBody(), SyncStateSet.class); return new Pair<>(header, stateSet); diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java index 9c6e05e3cb..cfe485c683 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java @@ -2443,7 +2443,7 @@ public class AdminBrokerProcessor implements NettyRequestProcessor { private RemotingCommand notifyBrokerRoleChanged(ChannelHandlerContext ctx, RemotingCommand request) throws RemotingCommandException { - NotifyBrokerRoleChangedRequestHeader requestHeader = request.decodeCommandCustomHeader(NotifyBrokerRoleChangedRequestHeader.class); + NotifyBrokerRoleChangedRequestHeader requestHeader = (NotifyBrokerRoleChangedRequestHeader) request.decodeCommandCustomHeader(NotifyBrokerRoleChangedRequestHeader.class); RemotingCommand response = RemotingCommand.createResponseCommand(null); diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java index ef63aad92d..a336ff5a2d 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java @@ -2889,7 +2889,7 @@ public class MQClientAPIImpl implements NameServerUpdateCallback { final RemotingCommand response = this.remotingClient.invokeSync(controllerAddress, request, 3000); assert response != null; if (response.getCode() == SUCCESS) { - return response.decodeCommandCustomHeader(GetMetaDataResponseHeader.class); + return (GetMetaDataResponseHeader) response.decodeCommandCustomHeader(GetMetaDataResponseHeader.class); } throw new MQBrokerException(response.getCode(), response.getRemark()); } diff --git a/common/src/main/java/org/apache/rocketmq/common/utils/FastJsonSerializer.java b/common/src/main/java/org/apache/rocketmq/common/utils/FastJsonSerializer.java index 805558d1d8..600054b40b 100644 --- a/common/src/main/java/org/apache/rocketmq/common/utils/FastJsonSerializer.java +++ b/common/src/main/java/org/apache/rocketmq/common/utils/FastJsonSerializer.java @@ -40,7 +40,7 @@ public class FastJsonSerializer implements Serializer { return new byte[0]; } else { try { - return JSON.toJSONBytesWithFastJsonConfig(this.fastJsonConfig.getCharset(), t, this.fastJsonConfig.getSerializeConfig(), this.fastJsonConfig.getSerializeFilters(), this.fastJsonConfig.getDateFormat(), JSON.DEFAULT_GENERATE_FEATURE, this.fastJsonConfig.getSerializerFeatures()); + return JSON.toJSONBytes(this.fastJsonConfig.getCharset(), t, this.fastJsonConfig.getSerializeConfig(), this.fastJsonConfig.getSerializeFilters(), this.fastJsonConfig.getDateFormat(), JSON.DEFAULT_GENERATE_FEATURE, this.fastJsonConfig.getSerializerFeatures()); } catch (Exception var3) { throw new SerializationException("Could not serialize: " + var3.getMessage(), var3); } diff --git a/controller/pom.xml b/controller/pom.xml index acddd82414..523475571d 100644 --- a/controller/pom.xml +++ b/controller/pom.xml @@ -57,10 +57,6 @@ ch.qos.logback logback-classic - - ch.qos.logback - logback-core - org.slf4j slf4j-api diff --git a/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/ReplicasInfoManager.java b/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/ReplicasInfoManager.java index 90ba39332c..be53e2d876 100644 --- a/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/ReplicasInfoManager.java +++ b/controller/src/main/java/org/apache/rocketmq/controller/impl/manager/ReplicasInfoManager.java @@ -168,7 +168,7 @@ public class ReplicasInfoManager { // Try elect a master in syncStateSet if (syncStateSet.size() > 1) { - boolean electSuccess = tryElectMaster(result, brokerName, syncStateSet, (candidate) -> + boolean electSuccess = tryElectMaster(result, brokerName, syncStateSet, candidate -> !candidate.equals(syncStateInfo.getMasterAddress()) && brokerAlivePredicate.test(brokerInfo.getClusterName(), candidate)); if (electSuccess) { return result; @@ -177,7 +177,7 @@ public class ReplicasInfoManager { // Try elect a master in lagging replicas if enableElectUncleanMaster = true if (controllerConfig.isEnableElectUncleanMaster()) { - boolean electSuccess = tryElectMaster(result, brokerName, brokerInfo.getAllBroker(), (candidate) -> + boolean electSuccess = tryElectMaster(result, brokerName, brokerInfo.getAllBroker(), candidate -> !candidate.equals(syncStateInfo.getMasterAddress()) && brokerAlivePredicate.test(brokerInfo.getClusterName(), candidate)); if (electSuccess) { return result; @@ -227,14 +227,15 @@ public class ReplicasInfoManager { final BrokerMemberGroup group = new BrokerMemberGroup(brokerInfo.getClusterName(), brokerName); final HashMap brokerIdTable = brokerInfo.getBrokerIdTable(); final HashMap memberGroup = new HashMap<>(); - brokerIdTable.forEach((addr, id)->memberGroup.put(id, addr)); + brokerIdTable.forEach((addr, id) -> memberGroup.put(id, addr)); group.setBrokerAddrs(memberGroup); return group; } return null; } - public ControllerResult registerBroker(final RegisterBrokerToControllerRequestHeader request) { + public ControllerResult registerBroker( + final RegisterBrokerToControllerRequestHeader request) { final String brokerName = request.getBrokerName(); final String brokerAddress = request.getBrokerAddress(); final ControllerResult result = new ControllerResult<>(new RegisterBrokerToControllerResponseHeader()); diff --git a/controller/src/main/java/org/apache/rocketmq/controller/processor/ControllerRequestProcessor.java b/controller/src/main/java/org/apache/rocketmq/controller/processor/ControllerRequestProcessor.java index b317e83837..765e67d199 100644 --- a/controller/src/main/java/org/apache/rocketmq/controller/processor/ControllerRequestProcessor.java +++ b/controller/src/main/java/org/apache/rocketmq/controller/processor/ControllerRequestProcessor.java @@ -71,7 +71,7 @@ public class ControllerRequestProcessor implements NettyRequestProcessor { } switch (request.getCode()) { case CONTROLLER_ALTER_SYNC_STATE_SET: { - final AlterSyncStateSetRequestHeader controllerRequest = request.decodeCommandCustomHeader(AlterSyncStateSetRequestHeader.class); + final AlterSyncStateSetRequestHeader controllerRequest = (AlterSyncStateSetRequestHeader) request.decodeCommandCustomHeader(AlterSyncStateSetRequestHeader.class); final SyncStateSet syncStateSet = RemotingSerializable.decode(request.getBody(), SyncStateSet.class); final CompletableFuture future = this.controller.alterSyncStateSet(controllerRequest, syncStateSet); if (future != null) { @@ -80,7 +80,7 @@ public class ControllerRequestProcessor implements NettyRequestProcessor { break; } case CONTROLLER_ELECT_MASTER: { - final ElectMasterRequestHeader controllerRequest = request.decodeCommandCustomHeader(ElectMasterRequestHeader.class); + final ElectMasterRequestHeader controllerRequest = (ElectMasterRequestHeader) request.decodeCommandCustomHeader(ElectMasterRequestHeader.class); final CompletableFuture future = this.controller.electMaster(controllerRequest); if (future != null) { return future.get(WAIT_TIMEOUT_OUT, TimeUnit.SECONDS); @@ -88,7 +88,7 @@ public class ControllerRequestProcessor implements NettyRequestProcessor { break; } case CONTROLLER_REGISTER_BROKER: { - final RegisterBrokerToControllerRequestHeader controllerRequest = request.decodeCommandCustomHeader(RegisterBrokerToControllerRequestHeader.class); + final RegisterBrokerToControllerRequestHeader controllerRequest = (RegisterBrokerToControllerRequestHeader) request.decodeCommandCustomHeader(RegisterBrokerToControllerRequestHeader.class); final CompletableFuture future = this.controller.registerBroker(controllerRequest); if (future != null) { final RemotingCommand response = future.get(WAIT_TIMEOUT_OUT, TimeUnit.SECONDS); @@ -102,7 +102,7 @@ public class ControllerRequestProcessor implements NettyRequestProcessor { break; } case CONTROLLER_GET_REPLICA_INFO: { - final GetReplicaInfoRequestHeader controllerRequest = request.decodeCommandCustomHeader(GetReplicaInfoRequestHeader.class); + final GetReplicaInfoRequestHeader controllerRequest = (GetReplicaInfoRequestHeader) request.decodeCommandCustomHeader(GetReplicaInfoRequestHeader.class); final CompletableFuture future = this.controller.getReplicaInfo(controllerRequest); if (future != null) { return future.get(WAIT_TIMEOUT_OUT, TimeUnit.SECONDS); @@ -113,7 +113,7 @@ public class ControllerRequestProcessor implements NettyRequestProcessor { return this.controller.getControllerMetadata(); } case BROKER_HEARTBEAT: { - final BrokerHeartbeatRequestHeader requestHeader = request.decodeCommandCustomHeader(BrokerHeartbeatRequestHeader.class); + final BrokerHeartbeatRequestHeader requestHeader = (BrokerHeartbeatRequestHeader) request.decodeCommandCustomHeader(BrokerHeartbeatRequestHeader.class); this.heartbeatManager.onBrokerHeartbeat(requestHeader.getClusterName(), requestHeader.getBrokerAddr()); return RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, "Heart beat success"); } diff --git a/controller/src/test/java/org/apache/rocketmq/controller/impl/controller/ControllerManagerTest.java b/controller/src/test/java/org/apache/rocketmq/controller/impl/controller/ControllerManagerTest.java index c9623f97a6..83a936c494 100644 --- a/controller/src/test/java/org/apache/rocketmq/controller/impl/controller/ControllerManagerTest.java +++ b/controller/src/test/java/org/apache/rocketmq/controller/impl/controller/ControllerManagerTest.java @@ -127,7 +127,7 @@ public class ControllerManagerTest { assert response != null; switch (response.getCode()) { case SUCCESS: { - return response.decodeCommandCustomHeader(RegisterBrokerToControllerResponseHeader.class); + return (RegisterBrokerToControllerResponseHeader) response.decodeCommandCustomHeader(RegisterBrokerToControllerResponseHeader.class); } case CONTROLLER_NOT_LEADER: { throw new MQBrokerException(response.getCode(), "Controller leader was changed"); @@ -175,7 +175,7 @@ public class ControllerManagerTest { final GetReplicaInfoRequestHeader requestHeader = new GetReplicaInfoRequestHeader("broker1"); final RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.CONTROLLER_GET_REPLICA_INFO, requestHeader); final RemotingCommand response = this.remotingClient1.invokeSync(leaderAddr, request, 3000); - final GetReplicaInfoResponseHeader responseHeader = response.decodeCommandCustomHeader(GetReplicaInfoResponseHeader.class); + final GetReplicaInfoResponseHeader responseHeader = (GetReplicaInfoResponseHeader) response.decodeCommandCustomHeader(GetReplicaInfoResponseHeader.class); assertEquals(responseHeader.getMasterAddress(), "127.0.0.1:8001"); executor.shutdown(); diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAService.java b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAService.java index 1c0a471ff9..32956dab2a 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAService.java +++ b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAService.java @@ -228,7 +228,7 @@ public class AutoSwitchHAService extends DefaultHAService { } public void updateConnectionLastCaughtUpTime(final String slaveAddress, final long lastCaughtUpTimeMs) { - long prevTime = this.connectionCaughtUpTimeTable.computeIfAbsent(slaveAddress, (k) -> 0L); + long prevTime = this.connectionCaughtUpTimeTable.computeIfAbsent(slaveAddress, k -> 0L); this.connectionCaughtUpTimeTable.put(slaveAddress, Math.max(prevTime, lastCaughtUpTimeMs)); } diff --git a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/EpochFileCache.java b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/EpochFileCache.java index f2b8b7d85c..9c310f8049 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/EpochFileCache.java +++ b/store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/EpochFileCache.java @@ -34,8 +34,7 @@ import org.apache.rocketmq.logging.InternalLogger; import org.apache.rocketmq.logging.InternalLoggerFactory; /** - * Cache for epochFile. - * Mapping (Epoch -> StartOffset) + * Cache for epochFile. Mapping (Epoch -> StartOffset) */ public class EpochFileCache { private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.STORE_LOGGER_NAME); @@ -249,7 +248,7 @@ public class EpochFileCache { * Remove epochEntries with epoch >= truncateEpoch. */ public void truncateSuffixByEpoch(final int truncateEpoch) { - Predicate predict = (entry) -> entry.getEpoch() >= truncateEpoch; + Predicate predict = entry -> entry.getEpoch() >= truncateEpoch; doTruncateSuffix(predict); } @@ -257,7 +256,7 @@ public class EpochFileCache { * Remove epochEntries with startOffset >= truncateOffset. */ public void truncateSuffixByOffset(final long truncateOffset) { - Predicate predict = (entry) -> entry.getStartOffset() >= truncateOffset; + Predicate predict = entry -> entry.getStartOffset() >= truncateOffset; doTruncateSuffix(predict); } @@ -279,7 +278,7 @@ public class EpochFileCache { * Remove epochEntries with endOffset <= truncateOffset. */ public void truncatePrefixByOffset(final long truncateOffset) { - Predicate predict = (entry) -> entry.getEndOffset() <= truncateOffset; + Predicate predict = entry -> entry.getEndOffset() <= truncateOffset; this.writeLock.lock(); try { this.epochMap.entrySet().removeIf(entry -> predict.test(entry.getValue()));