From 73618554f23f2c6792f00f5fb0104b4ebc7005c5 Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Wed, 1 Jun 2022 17:38:18 +0800 Subject: [PATCH] [ISSUE #3949] for code style --- .../rocketmq/common/constant/LoggerName.java | 1 + .../rocketmq/common/logger/ProxyLogger.java | 26 ----------------- .../common/logger/WatermarkLogger.java | 28 ------------------- .../common/thread/ThreadPoolMonitor.java | 27 ++++++++++++------ distribution/conf/logback_proxy.xml | 20 ++++++------- .../apache/rocketmq/proxy/ProxyStartup.java | 8 ++++-- .../proxy/common/ParameterConverter.java | 24 ---------------- .../proxy/common/utils/FilterUtils.java | 2 +- .../proxy/common/utils/FutureUtils.java | 3 +- .../rocketmq/proxy/config/ProxyConfig.java | 9 ++++++ .../interceptor/InterceptorConstants.java | 1 - .../grpc/v2/DefaultGrpcMessingActivity.java | 3 +- .../proxy/grpc/v2/GrpcMessingActivity.java | 12 +++++--- .../grpc/v2/channel/GrpcChannelManager.java | 4 +-- .../grpc/v2/channel/GrpcClientChannel.java | 3 +- .../proxy/grpc/v2/common/ResponseBuilder.java | 2 +- .../consumer/PopMessageResultFilterImpl.java | 3 +- .../v2/consumer/ReceiveMessageActivity.java | 3 +- .../proxy/processor/MessagingProcessor.java | 2 +- .../validator/TopicMessageTypeValidator.java | 3 +- .../proxy/service/ClusterServiceManager.java | 2 +- .../proxy/service/LocalServiceManager.java | 6 ++-- .../proxy/service/ServiceManager.java | 2 +- .../proxy/service/channel/SimpleChannel.java | 7 +++-- .../message/ClusterMessageService.java | 7 +++-- .../service/message/LocalMessageService.java | 12 +++++--- .../service/mqclient/MQClientAPIExt.java | 3 +- .../ProxyClientRemotingProcessor.java | 3 +- .../relay/ClusterProxyRelayService.java | 3 +- .../proxy/service/relay/ProxyChannel.java | 6 ++-- .../service/route/LocalTopicRouteService.java | 3 +- .../service/route/MessageQueueSelector.java | 3 +- .../proxy/service/route/MessageQueueView.java | 3 +- .../service/route/SelectableMessageQueue.java | 3 +- .../service/route/TopicRouteService.java | 6 ++-- .../ClusterTransactionService.java | 5 ++-- .../transaction/LocalTransactionService.java | 15 +++++++--- .../service/transaction/TransactionId.java | 6 ++-- .../proxy/grpc/v2/BaseActivityTest.java | 1 - .../grpc/v2/GrpcMessagingApplicationTest.java | 4 +-- .../grpc/v2/client/ClientActivityTest.java | 21 +++++++++----- .../common/GrpcClientSettingsManagerTest.java | 2 +- .../v2/consumer/AckMessageActivityTest.java | 3 +- .../ChangeInvisibleDurationActivityTest.java | 2 +- .../consumer/ReceiveMessageActivityTest.java | 1 - .../ForwardMessageToDLQActivityTest.java | 2 +- .../grpc/v2/route/RouteActivityTest.java | 28 +++++++++---------- .../EndTransactionActivityTest.java | 2 +- .../processor/ConsumerProcessorTest.java | 2 +- .../processor/ProducerProcessorTest.java | 2 +- .../metadata/ClusterMetadataServiceTest.java | 3 -- .../service/mqclient/MQClientAPIExtTest.java | 1 - .../proxy/service/relay/ProxyChannelTest.java | 11 +++++--- .../route/ClusterTopicRouteServiceTest.java | 5 ++-- .../ClusterTransactionServiceTest.java | 3 -- .../rmq-proxy-home/conf/logback_proxy.xml | 22 +++++++-------- 56 files changed, 189 insertions(+), 205 deletions(-) delete mode 100644 common/src/main/java/org/apache/rocketmq/common/logger/ProxyLogger.java delete mode 100644 common/src/main/java/org/apache/rocketmq/common/logger/WatermarkLogger.java delete mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/common/ParameterConverter.java diff --git a/common/src/main/java/org/apache/rocketmq/common/constant/LoggerName.java b/common/src/main/java/org/apache/rocketmq/common/constant/LoggerName.java index 244c5de772..b8ba059ea3 100644 --- a/common/src/main/java/org/apache/rocketmq/common/constant/LoggerName.java +++ b/common/src/main/java/org/apache/rocketmq/common/constant/LoggerName.java @@ -45,4 +45,5 @@ public class LoggerName { public static final String FAILOVER_LOGGER_NAME = "RocketmqFailover"; public static final String STDOUT_LOGGER_NAME = "STDOUT"; public static final String PROXY_LOGGER_NAME = "RocketmqProxy"; + public static final String PROXY_WATER_MARK_LOGGER_NAME = "RocketmqProxyWatermark"; } diff --git a/common/src/main/java/org/apache/rocketmq/common/logger/ProxyLogger.java b/common/src/main/java/org/apache/rocketmq/common/logger/ProxyLogger.java deleted file mode 100644 index 1e51ed9d8d..0000000000 --- a/common/src/main/java/org/apache/rocketmq/common/logger/ProxyLogger.java +++ /dev/null @@ -1,26 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.rocketmq.common.logger; - -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -public class ProxyLogger { - public static final Logger LOG_WATER_MARK = LoggerFactory.getLogger("Watermark"); - public static final Logger LOG_JSTACK = LoggerFactory.getLogger("Jstack"); -} diff --git a/common/src/main/java/org/apache/rocketmq/common/logger/WatermarkLogger.java b/common/src/main/java/org/apache/rocketmq/common/logger/WatermarkLogger.java deleted file mode 100644 index 9e57d02c8f..0000000000 --- a/common/src/main/java/org/apache/rocketmq/common/logger/WatermarkLogger.java +++ /dev/null @@ -1,28 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.apache.rocketmq.common.logger; - -import org.slf4j.Logger; - -public class WatermarkLogger { - private static final Logger LOG_MSG_TRACE = ProxyLogger.LOG_WATER_MARK; - - public static void info(String name, String k, double v) { - LOG_MSG_TRACE.info("\t{}\t{}\t{}", name, k, v); - } -} diff --git a/common/src/main/java/org/apache/rocketmq/common/thread/ThreadPoolMonitor.java b/common/src/main/java/org/apache/rocketmq/common/thread/ThreadPoolMonitor.java index 6db0b2f15a..e5bb6a394c 100644 --- a/common/src/main/java/org/apache/rocketmq/common/thread/ThreadPoolMonitor.java +++ b/common/src/main/java/org/apache/rocketmq/common/thread/ThreadPoolMonitor.java @@ -28,22 +28,30 @@ import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import org.apache.rocketmq.common.UtilAll; -import org.apache.rocketmq.common.logger.ProxyLogger; -import org.apache.rocketmq.common.logger.WatermarkLogger; +import org.apache.rocketmq.logging.InternalLogger; +import org.apache.rocketmq.logging.InternalLoggerFactory; public class ThreadPoolMonitor { + private static InternalLogger jstackLogger = InternalLoggerFactory.getLogger(ThreadPoolMonitor.class); + private static InternalLogger waterMarkLogger = InternalLoggerFactory.getLogger(ThreadPoolMonitor.class); + private static final List MONITOR_EXECUTOR = new CopyOnWriteArrayList<>(); private static final ScheduledExecutorService MONITOR_SCHEDULED = Executors.newSingleThreadScheduledExecutor( new ThreadFactoryBuilder().setNameFormat("ThreadPoolMonitor-%d").build() ); + private static volatile long threadPoolStatusPeriodTime = TimeUnit.SECONDS.toMillis(3); private static volatile boolean enablePrintJstack = true; - private static volatile long jstackPeriodTIme = 60000; + private static volatile long jstackPeriodTime = 60000; private static volatile long jstackTime = System.currentTimeMillis(); - public static void config(boolean enablePrintJstack, long jstackPeriodTime) { + public static void config(InternalLogger jstackLoggerConfig, InternalLogger waterMarkLoggerConfig, + boolean enablePrintJstack, long jstackPeriodTimeConfig, long threadPoolStatusPeriodTimeConfig) { + jstackLogger = jstackLoggerConfig; + waterMarkLogger = waterMarkLoggerConfig; + threadPoolStatusPeriodTime = threadPoolStatusPeriodTimeConfig; ThreadPoolMonitor.enablePrintJstack = enablePrintJstack; - jstackPeriodTIme = jstackPeriodTime; + jstackPeriodTime = jstackPeriodTimeConfig; } public static ThreadPoolExecutor createAndMonitor(int corePoolSize, @@ -97,15 +105,15 @@ public class ThreadPoolMonitor { List monitors = threadPoolWrapper.getStatusPrinters(); for (ThreadPoolStatusMonitor monitor : monitors) { double value = monitor.value(threadPoolWrapper.getThreadPoolExecutor()); - WatermarkLogger.info(threadPoolWrapper.getName(), + waterMarkLogger.info("\t{}\t{}\t{}", threadPoolWrapper.getName(), monitor.describe(), value); if (enablePrintJstack) { if (monitor.needPrintJstack(threadPoolWrapper.getThreadPoolExecutor(), value) && - System.currentTimeMillis() - jstackTime > jstackPeriodTIme) { + System.currentTimeMillis() - jstackTime > jstackPeriodTime) { jstackTime = System.currentTimeMillis(); - ProxyLogger.LOG_JSTACK.warn("jstack start \n " + UtilAll.jstack()); + jstackLogger.warn("jstack start\n{}", UtilAll.jstack()); } } } @@ -113,7 +121,8 @@ public class ThreadPoolMonitor { } public static void init() { - MONITOR_SCHEDULED.scheduleAtFixedRate(ThreadPoolMonitor::logThreadPoolStatus, 20, 1, TimeUnit.SECONDS); + MONITOR_SCHEDULED.scheduleAtFixedRate(ThreadPoolMonitor::logThreadPoolStatus, 20, + threadPoolStatusPeriodTime, TimeUnit.MILLISECONDS); } public static void shutdown() { diff --git a/distribution/conf/logback_proxy.xml b/distribution/conf/logback_proxy.xml index 8d0458ebf0..ad862d53c7 100644 --- a/distribution/conf/logback_proxy.xml +++ b/distribution/conf/logback_proxy.xml @@ -25,7 +25,7 @@ ${user.home}/logs/rocketmqlogs/otherdays/proxy.%i.log.gz 1 - 20 + 10 128MB @@ -39,25 +39,25 @@ - - ${user.home}/logs/rocketmqlogs/grpc.log + ${user.home}/logs/rocketmqlogs/proxy_watermark.log true - ${user.home}/logs/rocketmqlogs/otherdays/grpc.%i.log.gz + ${user.home}/logs/rocketmqlogs/otherdays/proxy_watermark.%i.log.gz 1 - 20 + 10 128MB - %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n + %d{yyy-MM-dd HH:mm:ss,GMT+8}%m%n UTF-8 - - + + @@ -408,9 +408,9 @@ - + - + diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java b/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java index d198830e30..3eb3c26850 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java @@ -149,9 +149,13 @@ public class ProxyStartup { } public static void initThreadPoolMonitor() { - ThreadPoolMonitor.init(); ProxyConfig config = ConfigurationManager.getProxyConfig(); - ThreadPoolMonitor.config(config.isEnablePrintJstack(), config.getPrintJstackInMillis()); + ThreadPoolMonitor.config( + InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME), + InternalLoggerFactory.getLogger(LoggerName.PROXY_WATER_MARK_LOGGER_NAME), + config.isEnablePrintJstack(), config.getPrintJstackInMillis(), + config.getPrintThreadPoolStatusInMillis()); + ThreadPoolMonitor.init(); } public static void initLogger() throws JoranException { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/common/ParameterConverter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/common/ParameterConverter.java deleted file mode 100644 index 2632d4340f..0000000000 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/common/ParameterConverter.java +++ /dev/null @@ -1,24 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.apache.rocketmq.proxy.common; - -import io.grpc.Context; - -@FunctionalInterface -public interface ParameterConverter { - R convert(Context ctx, T parameter) throws Throwable; -} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/FilterUtils.java b/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/FilterUtils.java index 2c9b663a9a..23eb1e1536 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/FilterUtils.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/FilterUtils.java @@ -24,7 +24,7 @@ public class FilterUtils { * Whether the message's tag matches consumerGroup's SubscriptionData * * @param tagsSet, tagSet in {@link SubscriptionData}, tagSet empty means SubscriptionData.SUB_ALL(*) - * @param tags, message's tags, null means not tag attached to the message. + * @param tags, message's tags, null means not tag attached to the message. */ public static boolean isTagMatched(Set tagsSet, String tags) { if (tagsSet.isEmpty()) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/FutureUtils.java b/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/FutureUtils.java index 3e3c5623ee..2e194a8cbe 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/FutureUtils.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/FutureUtils.java @@ -22,7 +22,8 @@ import java.util.concurrent.ExecutorService; public class FutureUtils { - public static CompletableFuture appendNextFuture(CompletableFuture future, CompletableFuture nextFuture, ExecutorService executor) { + public static CompletableFuture appendNextFuture(CompletableFuture future, + CompletableFuture nextFuture, ExecutorService executor) { future.whenCompleteAsync((t, throwable) -> { if (throwable != null) { nextFuture.completeExceptionally(throwable); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java b/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java index 6bf4ecf653..49e7b4ae80 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/config/ProxyConfig.java @@ -31,6 +31,7 @@ public class ProxyConfig { */ private boolean enablePrintJstack = true; private long printJstackInMillis = Duration.ofSeconds(60).toMillis(); + private long printThreadPoolStatusInMillis = Duration.ofSeconds(3).toMillis(); private String nameSrvAddr = ""; private String nameSrvDomain = ""; @@ -124,6 +125,14 @@ public class ProxyConfig { this.printJstackInMillis = printJstackInMillis; } + public long getPrintThreadPoolStatusInMillis() { + return printThreadPoolStatusInMillis; + } + + public void setPrintThreadPoolStatusInMillis(long printThreadPoolStatusInMillis) { + this.printThreadPoolStatusInMillis = printThreadPoolStatusInMillis; + } + public String getNameSrvAddr() { return nameSrvAddr; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/InterceptorConstants.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/InterceptorConstants.java index 62614b3a5e..c8aa39959e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/InterceptorConstants.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/InterceptorConstants.java @@ -35,7 +35,6 @@ public class InterceptorConstants { public static final Metadata.Key LOCAL_ADDRESS = Metadata.Key.of("rpc-local-address", Metadata.ASCII_STRING_MARSHALLER); - public static final Metadata.Key AUTHORIZATION = Metadata.Key.of("authorization", Metadata.ASCII_STRING_MARSHALLER); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/DefaultGrpcMessingActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/DefaultGrpcMessingActivity.java index 50d5f180a0..1422c01e78 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/DefaultGrpcMessingActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/DefaultGrpcMessingActivity.java @@ -103,7 +103,8 @@ public class DefaultGrpcMessingActivity extends AbstractStartAndShutdown impleme } @Override - public void receiveMessage(Context ctx, ReceiveMessageRequest request, StreamObserver responseObserver) { + public void receiveMessage(Context ctx, ReceiveMessageRequest request, + StreamObserver responseObserver) { this.receiveMessageActivity.receiveMessage(ctx, request, responseObserver); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessingActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessingActivity.java index 68337cba1c..796d5f57af 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessingActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessingActivity.java @@ -53,17 +53,21 @@ public interface GrpcMessingActivity extends StartAndShutdown { CompletableFuture queryAssignment(Context ctx, QueryAssignmentRequest request); - void receiveMessage(Context ctx, ReceiveMessageRequest request, StreamObserver responseObserver); + void receiveMessage(Context ctx, ReceiveMessageRequest request, + StreamObserver responseObserver); CompletableFuture ackMessage(Context ctx, AckMessageRequest request); - CompletableFuture forwardMessageToDeadLetterQueue(Context ctx, ForwardMessageToDeadLetterQueueRequest request); + CompletableFuture forwardMessageToDeadLetterQueue(Context ctx, + ForwardMessageToDeadLetterQueueRequest request); CompletableFuture endTransaction(Context ctx, EndTransactionRequest request); - CompletableFuture notifyClientTermination(Context ctx, NotifyClientTerminationRequest request); + CompletableFuture notifyClientTermination(Context ctx, + NotifyClientTerminationRequest request); - CompletableFuture changeInvisibleDuration(Context ctx, ChangeInvisibleDurationRequest request); + CompletableFuture changeInvisibleDuration(Context ctx, + ChangeInvisibleDurationRequest request); StreamObserver telemetry(Context ctx, StreamObserver responseObserver); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcChannelManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcChannelManager.java index fd07618c45..d063b6524b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcChannelManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcChannelManager.java @@ -54,7 +54,7 @@ public class GrpcChannelManager implements StartAndShutdown { protected void init() { this.scheduledExecutorService.scheduleAtFixedRate( this::scanExpireResultFuture, - 10, 10, TimeUnit.SECONDS + 10, 10, TimeUnit.SECONDS ); } @@ -77,7 +77,7 @@ public class GrpcChannelManager implements StartAndShutdown { return clientIdChannelMap.get(clientId); } - public GrpcClientChannel removeChannel(String group, String clientId) { + public GrpcClientChannel removeChannel(String group, String clientId) { AtomicReference channelRef = new AtomicReference<>(); this.groupClientIdChannelMap.computeIfPresent(group, (groupKey, clientIdMap) -> { channelRef.set(clientIdMap.remove(clientId)); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java index 2f25483041..a8492b140f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/channel/GrpcClientChannel.java @@ -50,7 +50,8 @@ public class GrpcClientChannel extends ProxyChannel { private final String group; private final String clientId; - public GrpcClientChannel(ProxyRelayService proxyRelayService, GrpcChannelManager grpcChannelManager, Context ctx, String group, String clientId) { + public GrpcClientChannel(ProxyRelayService proxyRelayService, GrpcChannelManager grpcChannelManager, Context ctx, + String group, String clientId) { super(proxyRelayService, null, new GrpcChannelId(group, clientId), InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.REMOTE_ADDRESS), InterceptorConstants.METADATA.get(ctx).get(InterceptorConstants.LOCAL_ADDRESS)); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/ResponseBuilder.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/ResponseBuilder.java index b68e1a77f9..412b01d34b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/ResponseBuilder.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common/ResponseBuilder.java @@ -62,7 +62,7 @@ public class ResponseBuilder { .setMessage(message) .build(); } - + public static Code buildCode(int remotingResponseCode) { switch (remotingResponseCode) { case ResponseCode.SUCCESS: diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/PopMessageResultFilterImpl.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/PopMessageResultFilterImpl.java index d411184149..3f9f2b2174 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/PopMessageResultFilterImpl.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/PopMessageResultFilterImpl.java @@ -32,7 +32,8 @@ public class PopMessageResultFilterImpl implements PopMessageResultFilter { } @Override - public FilterResult filterMessage(ProxyContext ctx, String consumerGroup, SubscriptionData subscriptionData, MessageExt messageExt) { + public FilterResult filterMessage(ProxyContext ctx, String consumerGroup, SubscriptionData subscriptionData, + MessageExt messageExt) { int maxAttempts = grpcClientSettingsManager.getClientSettings(ctx).getBackoffPolicy().getMaxAttempts(); if (!FilterUtils.isTagMatched(subscriptionData.getTagsSet(), messageExt.getTags())) { return FilterResult.NO_MATCH; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java index 0bd8bf565a..3c8b045a07 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java @@ -111,7 +111,8 @@ public class ReceiveMessageActivity extends AbstractMessingActivity { } } - protected ReceiveMessageResponseStreamWriter createWriter(ProxyContext ctx, StreamObserver responseObserver) { + protected ReceiveMessageResponseStreamWriter createWriter(ProxyContext ctx, + StreamObserver responseObserver) { return new ReceiveMessageResponseStreamWriter( this.messagingProcessor, responseObserver diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java index 1932c3f83b..07d28f8e46 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java @@ -57,7 +57,7 @@ public interface MessagingProcessor extends StartAndShutdown { ProxyContext ctx, List
requestHostAndPortList, String topicName - ) throws Exception; + ) throws Exception; default CompletableFuture> sendMessage( ProxyContext ctx, diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/validator/TopicMessageTypeValidator.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/validator/TopicMessageTypeValidator.java index 43eae1e314..137be90956 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/validator/TopicMessageTypeValidator.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/validator/TopicMessageTypeValidator.java @@ -22,8 +22,9 @@ import org.apache.rocketmq.common.attribute.TopicMessageType; public interface TopicMessageTypeValidator { /** * Will throw {@link org.apache.rocketmq.proxy.common.ProxyException} if validate failed. + * * @param topicMessageType Target topic - * @param messageType Message's type + * @param messageType Message's type */ void validate(TopicMessageType topicMessageType, TopicMessageType messageType); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/ClusterServiceManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/ClusterServiceManager.java index cd37acaf09..45a8d7cec6 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/ClusterServiceManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/ClusterServiceManager.java @@ -37,10 +37,10 @@ import org.apache.rocketmq.proxy.service.message.MessageService; import org.apache.rocketmq.proxy.service.metadata.ClusterMetadataService; import org.apache.rocketmq.proxy.service.metadata.MetadataService; import org.apache.rocketmq.proxy.service.mqclient.DoNothingClientRemotingProcessor; +import org.apache.rocketmq.proxy.service.mqclient.MQClientAPIFactory; import org.apache.rocketmq.proxy.service.mqclient.ProxyClientRemotingProcessor; import org.apache.rocketmq.proxy.service.relay.ClusterProxyRelayService; import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; -import org.apache.rocketmq.proxy.service.mqclient.MQClientAPIFactory; import org.apache.rocketmq.proxy.service.route.ClusterTopicRouteService; import org.apache.rocketmq.proxy.service.route.TopicRouteService; import org.apache.rocketmq.proxy.service.transaction.ClusterTransactionService; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/LocalServiceManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/LocalServiceManager.java index 42c9e249d2..c69b6773a0 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/LocalServiceManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/LocalServiceManager.java @@ -115,11 +115,13 @@ public class LocalServiceManager extends AbstractStartAndShutdown implements Ser } private class LocalServiceManagerStartAndShutdown implements StartAndShutdown { - @Override public void start() throws Exception { + @Override + public void start() throws Exception { LocalServiceManager.this.scheduledExecutorService.scheduleWithFixedDelay(channelManager::scanAndCleanChannels, 5, 5, TimeUnit.MINUTES); } - @Override public void shutdown() throws Exception { + @Override + public void shutdown() throws Exception { LocalServiceManager.this.scheduledExecutorService.shutdown(); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManager.java index 6a4f1cf371..563b567152 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManager.java @@ -20,8 +20,8 @@ import org.apache.rocketmq.broker.client.ConsumerManager; import org.apache.rocketmq.broker.client.ProducerManager; import org.apache.rocketmq.proxy.common.StartAndShutdown; import org.apache.rocketmq.proxy.service.message.MessageService; -import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; import org.apache.rocketmq.proxy.service.metadata.MetadataService; +import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; import org.apache.rocketmq.proxy.service.route.TopicRouteService; import org.apache.rocketmq.proxy.service.transaction.TransactionService; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/channel/SimpleChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/channel/SimpleChannel.java index 9f010526b8..4b700c5ed7 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/channel/SimpleChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/channel/SimpleChannel.java @@ -36,6 +36,7 @@ import org.apache.rocketmq.logging.InternalLoggerFactory; /** * SimpleChannel is used to handle writeAndFlush situation in processor + * * @see io.netty.channel.ChannelHandlerContext#writeAndFlush * @see io.netty.channel.Channel#writeAndFlush */ @@ -51,9 +52,9 @@ public class SimpleChannel extends AbstractChannel { /** * Creates a new instance. * - * @param parent the parent of this channel. {@code null} if there's no parent. - * @param remoteAddress Remote address - * @param localAddress Local address + * @param parent the parent of this channel. {@code null} if there's no parent. + * @param remoteAddress Remote address + * @param localAddress Local address */ public SimpleChannel(Channel parent, String remoteAddress, String localAddress) { super(parent); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/ClusterMessageService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/ClusterMessageService.java index 4fc8fdad9e..3238c69ba1 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/ClusterMessageService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/ClusterMessageService.java @@ -59,12 +59,12 @@ public class ClusterMessageService implements MessageService { future = this.mqClientAPIFactory.getClient().sendMessageAsync( messageQueue.getBrokerAddr(), messageQueue.getBrokerName(), msgList.get(0), requestHeader, timeoutMillis) - .thenApply(Lists::newArrayList); + .thenApply(Lists::newArrayList); } else { future = this.mqClientAPIFactory.getClient().sendMessageAsync( messageQueue.getBrokerAddr(), messageQueue.getBrokerName(), msgList, requestHeader, timeoutMillis) - .thenApply(Lists::newArrayList); + .thenApply(Lists::newArrayList); } return future; } @@ -81,7 +81,8 @@ public class ClusterMessageService implements MessageService { @Override public void endTransactionOneway(ProxyContext ctx, TransactionId transactionId, - EndTransactionRequestHeader requestHeader, long timeoutMillis) throws MQBrokerException, RemotingException, InterruptedException { + EndTransactionRequestHeader requestHeader, + long timeoutMillis) throws MQBrokerException, RemotingException, InterruptedException { this.mqClientAPIFactory.getClient().endTransactionOneway( this.resolveBrokerAddr(transactionId.getBrokerName()), requestHeader, diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/LocalMessageService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/LocalMessageService.java index 4e893c384f..1c079d7fb3 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/LocalMessageService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/message/LocalMessageService.java @@ -75,7 +75,8 @@ public class LocalMessageService implements MessageService { this.channelManager = channelManager; } - @Override public CompletableFuture> sendMessage(ProxyContext ctx, SelectableMessageQueue messageQueue, + @Override + public CompletableFuture> sendMessage(ProxyContext ctx, SelectableMessageQueue messageQueue, List msgList, SendMessageRequestHeader requestHeader, long timeoutMillis) { byte[] body; String messageId; @@ -162,7 +163,8 @@ public class LocalMessageService implements MessageService { return future; } - @Override public void endTransactionOneway(ProxyContext ctx, TransactionId transactionId, + @Override + public void endTransactionOneway(ProxyContext ctx, TransactionId transactionId, EndTransactionRequestHeader requestHeader, long timeoutMillis) { SimpleChannel channel = channelManager.createChannel(ctx); ChannelHandlerContext channelHandlerContext = channel.getChannelHandlerContext(); @@ -175,7 +177,8 @@ public class LocalMessageService implements MessageService { } } - @Override public CompletableFuture popMessage(ProxyContext ctx, SelectableMessageQueue messageQueue, + @Override + public CompletableFuture popMessage(ProxyContext ctx, SelectableMessageQueue messageQueue, PopMessageRequestHeader requestHeader, long timeoutMillis) { RemotingCommand request = LocalRemotingCommand.createRequestCommand(RequestCode.POP_MESSAGE, requestHeader); CompletableFuture future = new CompletableFuture<>(); @@ -315,7 +318,8 @@ public class LocalMessageService implements MessageService { }); } - @Override public CompletableFuture ackMessage(ProxyContext ctx, ReceiptHandle handle, String messageId, + @Override + public CompletableFuture ackMessage(ProxyContext ctx, ReceiptHandle handle, String messageId, AckMessageRequestHeader requestHeader, long timeoutMillis) { SimpleChannel channel = channelManager.createChannel(ctx); ChannelHandlerContext channelHandlerContext = channel.getChannelHandlerContext(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/mqclient/MQClientAPIExt.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/mqclient/MQClientAPIExt.java index 011500f8a2..1af6f3acb7 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/mqclient/MQClientAPIExt.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/mqclient/MQClientAPIExt.java @@ -431,7 +431,8 @@ public class MQClientAPIExt extends MQClientAPIImpl { return future; } - public CompletableFuture searchOffsetAsync(String brokerAddr, String topic, int queueId , long timestamp, long timeoutMillis) { + public CompletableFuture searchOffsetAsync(String brokerAddr, String topic, int queueId, long timestamp, + long timeoutMillis) { SearchOffsetRequestHeader requestHeader = new SearchOffsetRequestHeader(); requestHeader.setTopic(topic); requestHeader.setQueueId(queueId); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/mqclient/ProxyClientRemotingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/mqclient/ProxyClientRemotingProcessor.java index ca3edf3ef7..d932cd1596 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/mqclient/ProxyClientRemotingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/mqclient/ProxyClientRemotingProcessor.java @@ -46,7 +46,8 @@ public class ProxyClientRemotingProcessor extends ClientRemotingProcessor { } @Override - public RemotingCommand checkTransactionState(ChannelHandlerContext ctx, RemotingCommand request) throws RemotingCommandException { + public RemotingCommand checkTransactionState(ChannelHandlerContext ctx, + RemotingCommand request) throws RemotingCommandException { final ByteBuffer byteBuffer = ByteBuffer.wrap(request.getBody()); final MessageExt messageExt = MessageDecoder.decode(byteBuffer, true, false, false); if (messageExt != null) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ClusterProxyRelayService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ClusterProxyRelayService.java index 9b356e788b..fd7afaec9a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ClusterProxyRelayService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ClusterProxyRelayService.java @@ -36,7 +36,8 @@ public class ClusterProxyRelayService implements ProxyRelayService { return null; } - @Override public CompletableFuture> processConsumeMessageDirectly( + @Override + public CompletableFuture> processConsumeMessageDirectly( ProxyContext context, RemotingCommand command, ConsumeMessageDirectlyResultRequestHeader header) { return null; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ProxyChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ProxyChannel.java index 54bf8d0a7d..12a6ae6541 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ProxyChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/relay/ProxyChannel.java @@ -56,7 +56,8 @@ public abstract class ProxyChannel extends AbstractChannel { protected final ProxyRelayService proxyRelayService; - protected ProxyChannel(ProxyRelayService proxyRelayService, Channel parent, String remoteAddress, String localAddress) { + protected ProxyChannel(ProxyRelayService proxyRelayService, Channel parent, String remoteAddress, + String localAddress) { super(parent); this.proxyRelayService = proxyRelayService; this.remoteAddress = remoteAddress; @@ -65,7 +66,8 @@ public abstract class ProxyChannel extends AbstractChannel { this.localSocketAddress = RemotingUtil.string2SocketAddress(localAddress); } - protected ProxyChannel(ProxyRelayService proxyRelayService, Channel parent, ChannelId id, String remoteAddress, String localAddress) { + protected ProxyChannel(ProxyRelayService proxyRelayService, Channel parent, ChannelId id, String remoteAddress, + String localAddress) { super(parent, id); this.proxyRelayService = proxyRelayService; this.remoteAddress = remoteAddress; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/LocalTopicRouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/LocalTopicRouteService.java index 5da75cea08..0c83da60c2 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/LocalTopicRouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/LocalTopicRouteService.java @@ -56,7 +56,8 @@ public class LocalTopicRouteService extends TopicRouteService { } @Override - public ProxyTopicRouteData getTopicRouteForProxy(List
requestHostAndPortList, String topicName) throws Exception { + public ProxyTopicRouteData getTopicRouteForProxy(List
requestHostAndPortList, + String topicName) throws Exception { MessageQueueView messageQueueView = getAllMessageQueueView(topicName); TopicRouteData topicRouteData = messageQueueView.getTopicRouteData(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelector.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelector.java index 7c2f40a6e0..02c817f6f3 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelector.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueSelector.java @@ -193,7 +193,8 @@ public class MessageQueueSelector { return Objects.hash(queues, brokerActingQueues); } - @Override public String toString() { + @Override + public String toString() { return MoreObjects.toStringHelper(this) .add("queues", queues) .add("brokerActingQueues", brokerActingQueues) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueView.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueView.java index 06c9d14569..cdef39cc2d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueView.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/MessageQueueView.java @@ -53,7 +53,8 @@ public class MessageQueueView { return writeSelector; } - @Override public String toString() { + @Override + public String toString() { return MoreObjects.toStringHelper(this) .add("readSelector", readSelector) .add("writeSelector", writeSelector) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/SelectableMessageQueue.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/SelectableMessageQueue.java index 85f7434aed..99eccbedfc 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/SelectableMessageQueue.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/SelectableMessageQueue.java @@ -72,7 +72,8 @@ public class SelectableMessageQueue implements Comparable requestHostAndPortList, String topicName) throws Exception; + public abstract ProxyTopicRouteData getTopicRouteForProxy(List
requestHostAndPortList, + String topicName) throws Exception; public abstract String getBrokerAddr(String brokerName) throws Exception; - protected static MessageQueueView getCacheMessageQueueWrapper(LoadingCache topicCache, String key) throws Exception { + protected static MessageQueueView getCacheMessageQueueWrapper(LoadingCache topicCache, + String key) throws Exception { MessageQueueView res = topicCache.get(key); if (res.isEmptyCachedQueue()) { throw new MQClientException(ResponseCode.TOPIC_NOT_EXIST, diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/transaction/ClusterTransactionService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/transaction/ClusterTransactionService.java index 210f092203..93966d4f79 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/transaction/ClusterTransactionService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/transaction/ClusterTransactionService.java @@ -57,7 +57,8 @@ public class ClusterTransactionService implements StartAndShutdown, TransactionS private final Map/* cluster list */> groupClusterData = new ConcurrentHashMap<>(); private TxHeartbeatServiceThread txHeartbeatServiceThread; - public ClusterTransactionService(TopicRouteService topicRouteService, ProducerManager producerManager, RPCHook rpcHook, + public ClusterTransactionService(TopicRouteService topicRouteService, ProducerManager producerManager, + RPCHook rpcHook, MQClientAPIFactory mqClientAPIFactory) { this.topicRouteService = topicRouteService; this.mqClientAPIFactory = mqClientAPIFactory; @@ -186,7 +187,7 @@ public class ClusterTransactionService implements StartAndShutdown, TransactionS protected void sendHeartBeatToCluster(String clusterName, HeartbeatData heartbeatData) { try { - MessageQueueView messageQueue = this.topicRouteService.getAllMessageQueueView(clusterName); + MessageQueueView messageQueue = this.topicRouteService.getAllMessageQueueView(clusterName); List brokerDataList = messageQueue.getTopicRouteData().getBrokerDatas(); if (brokerDataList == null) { return; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/transaction/LocalTransactionService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/transaction/LocalTransactionService.java index c465520a98..fe0a1c0451 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/transaction/LocalTransactionService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/transaction/LocalTransactionService.java @@ -18,20 +18,27 @@ package org.apache.rocketmq.proxy.service.transaction; import java.util.List; +/** + * no need to implements, because the channel of producer will put into the broker's producerManager + */ public class LocalTransactionService implements TransactionService { - @Override public void addTransactionSubscription(String group, List topicList) { + @Override + public void addTransactionSubscription(String group, List topicList) { } - @Override public void addTransactionSubscription(String group, String topic) { + @Override + public void addTransactionSubscription(String group, String topic) { } - @Override public void replaceTransactionSubscription(String group, List topicList) { + @Override + public void replaceTransactionSubscription(String group, List topicList) { } - @Override public void unSubscribeAllTransactionTopic(String group) { + @Override + public void unSubscribeAllTransactionTopic(String group) { } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/transaction/TransactionId.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/transaction/TransactionId.java index 44b5056911..78b479634a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/transaction/TransactionId.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/transaction/TransactionId.java @@ -181,7 +181,8 @@ public class TransactionId { this.proxyTransactionId = proxyTransactionId; } - @Override public String toString() { + @Override + public String toString() { return MoreObjects.toStringHelper(this) .add("brokerName", brokerName) .add("brokerTransactionId", brokerTransactionId) @@ -230,7 +231,8 @@ public class TransactionId { return new TransactionId(brokerName, brokerTransactionId, commitLogOffset, tranStateTableOffset, proxyTransactionId); } - @Override public String toString() { + @Override + public String toString() { return MoreObjects.toStringHelper(this) .add("brokerName", brokerName) .add("brokerTransactionId", brokerTransactionId) diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/BaseActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/BaseActivityTest.java index c17cca7c2c..bb28e43aea 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/BaseActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/BaseActivityTest.java @@ -31,7 +31,6 @@ import org.apache.rocketmq.proxy.processor.MessagingProcessor; import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; import org.junit.Ignore; import org.junit.runner.RunWith; -import org.mockito.Mock; import org.mockito.junit.MockitoJUnitRunner; import static org.mockito.Mockito.mock; diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingApplicationTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingApplicationTest.java index 5fc223787c..79999164f4 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingApplicationTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingApplicationTest.java @@ -48,7 +48,7 @@ public class GrpcMessagingApplicationTest extends InitConfigAndLoggerTest { GrpcMessagingApplication grpcMessagingApplication; private static final String TOPIC = "topic"; - private static Endpoints GRPC_ENDPOINTS = Endpoints.newBuilder() + private static Endpoints grpcEndpoints = Endpoints.newBuilder() .setScheme(AddressScheme.IPv4) .addAddresses(Address.newBuilder().setHost("127.0.0.1").setPort(8080).build()) .addAddresses(Address.newBuilder().setHost("127.0.0.2").setPort(8080).build()) @@ -64,7 +64,7 @@ public class GrpcMessagingApplicationTest extends InitConfigAndLoggerTest { public void testQueryRoute() { CompletableFuture future = new CompletableFuture<>(); QueryRouteRequest request = QueryRouteRequest.newBuilder() - .setEndpoints(GRPC_ENDPOINTS) + .setEndpoints(grpcEndpoints) .setTopic(Resource.newBuilder().setName(TOPIC).build()) .build(); Mockito.when(grpcMessingActivity.queryRoute(Mockito.any(Context.class), Mockito.eq(request))) diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivityTest.java index b295cbb477..242d62ee5a 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/client/ClientActivityTest.java @@ -261,13 +261,16 @@ public class ClientActivityTest extends BaseActivityTest { when(grpcChannelManagerMock.getAndRemoveResponseFuture(anyString())).thenReturn((CompletableFuture) runningInfoFutureMock); Context context = createContext(); StreamObserver streamObserver = clientActivity.telemetry(context, new StreamObserver() { - @Override public void onNext(TelemetryCommand value) { + @Override + public void onNext(TelemetryCommand value) { } - @Override public void onError(Throwable t) { + @Override + public void onError(Throwable t) { } - @Override public void onCompleted() { + @Override + public void onCompleted() { } }); streamObserver.onNext(TelemetryCommand.newBuilder() @@ -290,13 +293,16 @@ public class ClientActivityTest extends BaseActivityTest { when(grpcChannelManagerMock.getAndRemoveResponseFuture(anyString())).thenReturn((CompletableFuture) resultFutureMock); Context context = createContext(); StreamObserver streamObserver = clientActivity.telemetry(context, new StreamObserver() { - @Override public void onNext(TelemetryCommand value) { + @Override + public void onNext(TelemetryCommand value) { } - @Override public void onError(Throwable t) { + @Override + public void onError(Throwable t) { } - @Override public void onCompleted() { + @Override + public void onCompleted() { } }); streamObserver.onNext(TelemetryCommand.newBuilder() @@ -326,7 +332,8 @@ public class ClientActivityTest extends BaseActivityTest { } - @Override public void onCompleted() { + @Override + public void onCompleted() { } }; diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManagerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManagerTest.java index bf89caab64..28a422df0f 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManagerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/common/GrpcClientSettingsManagerTest.java @@ -67,7 +67,7 @@ public class GrpcClientSettingsManagerTest extends BaseActivityTest { subscriptionGroupConfig.setRetryMaxTimes(3); subscriptionGroupConfig.getGroupRetryPolicy().setType(GroupRetryPolicyType.CUSTOMIZED); - subscriptionGroupConfig.getGroupRetryPolicy().setCustomizedRetryPolicy(new CustomizedRetryPolicy(new long[]{1000})); + subscriptionGroupConfig.getGroupRetryPolicy().setCustomizedRetryPolicy(new CustomizedRetryPolicy(new long[] {1000})); settings = this.grpcClientSettingsManager.getClientSettings(context); assertEquals(RetryPolicy.newBuilder() .setMaxAttempts(3) diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/AckMessageActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/AckMessageActivityTest.java index 6ca311d6ff..d4a34cbb0c 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/AckMessageActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/AckMessageActivityTest.java @@ -24,14 +24,13 @@ import apache.rocketmq.v2.Code; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.client.consumer.AckResult; import org.apache.rocketmq.client.consumer.AckStatus; -import org.apache.rocketmq.common.protocol.header.ExtraInfoUtil; import org.apache.rocketmq.proxy.common.ProxyException; import org.apache.rocketmq.proxy.common.ProxyExceptionCode; import org.apache.rocketmq.proxy.grpc.v2.BaseActivityTest; import org.junit.Before; import org.junit.Test; -import static org.junit.Assert.*; +import static org.junit.Assert.assertEquals; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ChangeInvisibleDurationActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ChangeInvisibleDurationActivityTest.java index cf1c8b7983..4d6655530f 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ChangeInvisibleDurationActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ChangeInvisibleDurationActivityTest.java @@ -31,7 +31,7 @@ import org.junit.Before; import org.junit.Test; import org.mockito.ArgumentCaptor; -import static org.junit.Assert.*; +import static org.junit.Assert.assertEquals; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.when; diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java index 412e5806c9..40a5ed4243 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java @@ -42,7 +42,6 @@ import org.apache.rocketmq.proxy.common.ProxyContext; import org.apache.rocketmq.proxy.grpc.v2.BaseActivityTest; import org.apache.rocketmq.proxy.service.route.MessageQueueView; import org.apache.rocketmq.proxy.service.route.SelectableMessageQueue; -import org.assertj.core.util.Lists; import org.junit.Before; import org.junit.Test; import org.mockito.ArgumentCaptor; diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/producer/ForwardMessageToDLQActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/producer/ForwardMessageToDLQActivityTest.java index c6153d7133..1bae776f9c 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/producer/ForwardMessageToDLQActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/producer/ForwardMessageToDLQActivityTest.java @@ -29,7 +29,7 @@ import org.apache.rocketmq.remoting.protocol.RemotingCommand; import org.junit.Before; import org.junit.Test; -import static org.junit.Assert.*; +import static org.junit.Assert.assertEquals; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.when; diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java index ebb821fd3b..583dbc995e 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java @@ -58,18 +58,18 @@ public class RouteActivityTest extends BaseActivityTest { private static final Resource GRPC_TOPIC = Resource.newBuilder() .setName(TOPIC) .build(); - private static Endpoints GRPC_ENDPOINTS = Endpoints.newBuilder() + private static Endpoints grpcEndpoints = Endpoints.newBuilder() .setScheme(AddressScheme.IPv4) .addAddresses(Address.newBuilder().setHost("127.0.0.1").setPort(8080).build()) .addAddresses(Address.newBuilder().setHost("127.0.0.2").setPort(8080).build()) .build(); - private static List ENDPOINTS_ADDRESS = new ArrayList<>(); + private static List addressArrayList = new ArrayList<>(); static { - ENDPOINTS_ADDRESS.add(new org.apache.rocketmq.proxy.common.Address( + addressArrayList.add(new org.apache.rocketmq.proxy.common.Address( org.apache.rocketmq.proxy.common.Address.AddressScheme.IPv4, HostAndPort.fromParts("127.0.0.1", 8080))); - ENDPOINTS_ADDRESS.add(new org.apache.rocketmq.proxy.common.Address( + addressArrayList.add(new org.apache.rocketmq.proxy.common.Address( org.apache.rocketmq.proxy.common.Address.AddressScheme.IPv4, HostAndPort.fromParts("127.0.0.2", 8080))); } @@ -89,16 +89,16 @@ public class RouteActivityTest extends BaseActivityTest { QueryRouteResponse response = this.routeActivity.queryRoute( createContext(), QueryRouteRequest.newBuilder() - .setEndpoints(GRPC_ENDPOINTS) + .setEndpoints(grpcEndpoints) .setTopic(Resource.newBuilder().setName(TOPIC).build()) .build() ).get(); - assertEquals(ENDPOINTS_ADDRESS, addressListCaptor.getValue()); + assertEquals(addressArrayList, addressListCaptor.getValue()); assertEquals(Code.OK, response.getStatus().getCode()); assertEquals(4, response.getMessageQueuesCount()); for (MessageQueue messageQueue : response.getMessageQueuesList()) { - assertEquals(GRPC_ENDPOINTS, messageQueue.getBroker().getEndpoints()); + assertEquals(grpcEndpoints, messageQueue.getBroker().getEndpoints()); assertEquals(Permission.READ_WRITE, messageQueue.getPermission()); } } @@ -111,7 +111,7 @@ public class RouteActivityTest extends BaseActivityTest { QueryRouteResponse response = this.routeActivity.queryRoute( createContext(), QueryRouteRequest.newBuilder() - .setEndpoints(GRPC_ENDPOINTS) + .setEndpoints(grpcEndpoints) .setTopic(GRPC_TOPIC) .build() ).get(); @@ -127,7 +127,7 @@ public class RouteActivityTest extends BaseActivityTest { QueryAssignmentResponse response = this.routeActivity.queryAssignment( createContext(), QueryAssignmentRequest.newBuilder() - .setEndpoints(GRPC_ENDPOINTS) + .setEndpoints(grpcEndpoints) .setTopic(GRPC_TOPIC) .build() ).get(); @@ -143,7 +143,7 @@ public class RouteActivityTest extends BaseActivityTest { QueryAssignmentResponse response = this.routeActivity.queryAssignment( createContext(), QueryAssignmentRequest.newBuilder() - .setEndpoints(GRPC_ENDPOINTS) + .setEndpoints(grpcEndpoints) .setTopic(GRPC_TOPIC) .build() ).get(); @@ -159,14 +159,14 @@ public class RouteActivityTest extends BaseActivityTest { QueryAssignmentResponse response = this.routeActivity.queryAssignment( createContext(), QueryAssignmentRequest.newBuilder() - .setEndpoints(GRPC_ENDPOINTS) + .setEndpoints(grpcEndpoints) .setTopic(GRPC_TOPIC) .build() ).get(); assertEquals(Code.OK, response.getStatus().getCode()); assertEquals(1, response.getAssignmentsCount()); - assertEquals(GRPC_ENDPOINTS, response.getAssignments(0).getMessageQueue().getBroker().getEndpoints()); + assertEquals(grpcEndpoints, response.getAssignments(0).getMessageQueue().getBroker().getEndpoints()); } private static ProxyTopicRouteData createProxyTopicRouteData(int r, int w, int p) { @@ -175,8 +175,8 @@ public class RouteActivityTest extends BaseActivityTest { ProxyTopicRouteData.ProxyBrokerData proxyBrokerData = new ProxyTopicRouteData.ProxyBrokerData(); proxyBrokerData.setCluster(CLUSTER); proxyBrokerData.setBrokerName(BROKER_NAME); - proxyBrokerData.getBrokerAddrs().put(0L, ENDPOINTS_ADDRESS); - proxyBrokerData.getBrokerAddrs().put(1L, ENDPOINTS_ADDRESS); + proxyBrokerData.getBrokerAddrs().put(0L, addressArrayList); + proxyBrokerData.getBrokerAddrs().put(1L, addressArrayList); proxyTopicRouteData.getBrokerDatas().add(proxyBrokerData); return proxyTopicRouteData; } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/transaction/EndTransactionActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/transaction/EndTransactionActivityTest.java index 6709ae05ff..a2444a6e88 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/transaction/EndTransactionActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/transaction/EndTransactionActivityTest.java @@ -91,7 +91,7 @@ public class EndTransactionActivityTest extends BaseActivityTest { @Parameterized.Parameters public static Collection parameters() { - Object[][] p = new Object[][]{ + Object[][] p = new Object[][] { {TransactionResolution.COMMIT, TransactionSource.SOURCE_CLIENT, TransactionStatus.COMMIT, false}, {TransactionResolution.ROLLBACK, TransactionSource.SOURCE_SERVER_CHECK, TransactionStatus.ROLLBACK, true}, {TransactionResolution.TRANSACTION_RESOLUTION_UNSPECIFIED, TransactionSource.SOURCE_SERVER_CHECK, TransactionStatus.UNKNOWN, true}, diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java index dbac6b4e99..35ab32a9e9 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ConsumerProcessorTest.java @@ -58,7 +58,7 @@ public class ConsumerProcessorTest extends BaseProcessorTest { private static final String CONSUMER_GROUP = "consumerGroup"; private static final String TOPIC = "topic"; - + private ConsumerProcessor consumerProcessor; @Before diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java index 7208c6da83..5b66fa8ed0 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java @@ -77,7 +77,7 @@ public class ProducerProcessorTest extends BaseProcessorTest { sendResult.setMsgId(msgId); ArgumentCaptor requestHeaderArgumentCaptor = ArgumentCaptor.forClass(SendMessageRequestHeader.class); when(this.messageService.sendMessage(any(), any(), any(), requestHeaderArgumentCaptor.capture(), anyLong())) - .thenReturn(CompletableFuture.completedFuture(Lists.newArrayList(sendResult))); + .thenReturn(CompletableFuture.completedFuture(Lists.newArrayList(sendResult))); List messageExtList = new ArrayList<>(); MessageExt messageExt = createMessageExt(TOPIC, "tag", 0, 0); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/metadata/ClusterMetadataServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/metadata/ClusterMetadataServiceTest.java index 7eeb72060e..2c0d3f8909 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/metadata/ClusterMetadataServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/metadata/ClusterMetadataServiceTest.java @@ -18,14 +18,11 @@ package org.apache.rocketmq.proxy.service.metadata; import java.util.HashMap; -import org.apache.rocketmq.client.exception.MQClientException; import org.apache.rocketmq.common.attribute.TopicMessageType; -import org.apache.rocketmq.common.protocol.ResponseCode; import org.apache.rocketmq.common.statictopic.TopicConfigAndQueueMapping; import org.apache.rocketmq.common.subscription.SubscriptionGroupConfig; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.service.BaseServiceTest; -import org.apache.rocketmq.proxy.service.route.MessageQueueView; import org.junit.Before; import org.junit.Test; diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/mqclient/MQClientAPIExtTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/mqclient/MQClientAPIExtTest.java index b78502e2c1..db3ef7bf8d 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/mqclient/MQClientAPIExtTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/mqclient/MQClientAPIExtTest.java @@ -313,7 +313,6 @@ public class MQClientAPIExtTest { assertEquals(offset, mqClientAPI.getMaxOffsetAsync(BROKER_ADDR, TOPIC, 0, TIMEOUT).get().longValue()); } - @Test public void testSearchOffsetAsync() throws Exception { long offset = ThreadLocalRandom.current().nextLong(); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/relay/ProxyChannelTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/relay/ProxyChannelTest.java index a6d6d60f14..fadb280e37 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/relay/ProxyChannelTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/relay/ProxyChannelTest.java @@ -18,7 +18,6 @@ package org.apache.rocketmq.proxy.service.relay; import io.netty.channel.Channel; -import java.net.SocketAddress; import java.nio.charset.StandardCharsets; import java.util.UUID; import java.util.concurrent.CompletableFuture; @@ -40,7 +39,9 @@ import org.mockito.ArgumentCaptor; import org.mockito.Mock; import org.mockito.junit.MockitoJUnitRunner; -import static org.junit.Assert.*; +import static org.junit.Assert.assertArrayEquals; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.when; @@ -57,11 +58,13 @@ public class ProxyChannelTest { super(proxyRelayService, parent, remoteAddress, localAddress); } - @Override public boolean isOpen() { + @Override + public boolean isOpen() { return false; } - @Override public boolean isActive() { + @Override + public boolean isActive() { return false; } } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/ClusterTopicRouteServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/ClusterTopicRouteServiceTest.java index 2d5f64de37..2a5d3189eb 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/ClusterTopicRouteServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/ClusterTopicRouteServiceTest.java @@ -18,7 +18,6 @@ package org.apache.rocketmq.proxy.service.route; import com.google.common.net.HostAndPort; -import java.util.ArrayList; import java.util.List; import org.apache.rocketmq.client.exception.MQClientException; import org.apache.rocketmq.common.MixAll; @@ -30,7 +29,9 @@ import org.junit.Before; import org.junit.Test; import static org.assertj.core.api.Assertions.catchThrowableOfType; -import static org.junit.Assert.*; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.when; diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/transaction/ClusterTransactionServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/transaction/ClusterTransactionServiceTest.java index 72fad8f13a..c9b3b17657 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/transaction/ClusterTransactionServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/transaction/ClusterTransactionServiceTest.java @@ -30,10 +30,8 @@ import org.apache.rocketmq.proxy.service.route.MessageQueueView; import org.assertj.core.util.Lists; import org.junit.Before; import org.junit.Test; -import org.junit.runner.RunWith; import org.mockito.ArgumentCaptor; import org.mockito.Mock; -import org.mockito.junit.MockitoJUnitRunner; import static org.awaitility.Awaitility.await; import static org.junit.Assert.assertEquals; @@ -43,7 +41,6 @@ import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.doNothing; import static org.mockito.Mockito.when; - public class ClusterTransactionServiceTest extends BaseServiceTest { @Mock diff --git a/proxy/src/test/resources/rmq-proxy-home/conf/logback_proxy.xml b/proxy/src/test/resources/rmq-proxy-home/conf/logback_proxy.xml index 8d0458ebf0..74829684e8 100644 --- a/proxy/src/test/resources/rmq-proxy-home/conf/logback_proxy.xml +++ b/proxy/src/test/resources/rmq-proxy-home/conf/logback_proxy.xml @@ -25,7 +25,7 @@ ${user.home}/logs/rocketmqlogs/otherdays/proxy.%i.log.gz 1 - 20 + 10 128MB @@ -39,25 +39,25 @@ - - ${user.home}/logs/rocketmqlogs/grpc.log + ${user.home}/logs/rocketmqlogs/proxy_watermark.log true - ${user.home}/logs/rocketmqlogs/otherdays/grpc.%i.log.gz + ${user.home}/logs/rocketmqlogs/otherdays/proxy_watermark.%i.log.gz 1 - 20 + 10 128MB - %d{yyy-MM-dd HH:mm:ss,GMT+8} %p %t - %m%n + %d{yyy-MM-dd HH:mm:ss,GMT+8}%m%n UTF-8 - - + + @@ -312,7 +312,7 @@ 20 + class="ch.qos.logback.core.rolling.SizeBasedTriggeringPolicy"> 128MB @@ -408,9 +408,9 @@ - + - +