This is an automated email from the ASF dual-hosted git repository.

mikexue pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-eventmesh.git


The following commit(s) were added to refs/heads/master by this push:
     new 6f6629e55 simplify code
     new 826aab961 Merge pull request #2754 from weihubeats/ConsumerManager
6f6629e55 is described below

commit 6f6629e55d3ee3b8bddf2760b4015709358c9440
Author: weihu <[email protected]>
AuthorDate: Sat Dec 31 15:50:19 2022 +0800

    simplify code
---
 .../protocol/grpc/consumer/ConsumerManager.java    | 26 ++++++++++------------
 1 file changed, 12 insertions(+), 14 deletions(-)

diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/consumer/ConsumerManager.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/consumer/ConsumerManager.java
index 7caf6b222..63456437d 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/consumer/ConsumerManager.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/consumer/ConsumerManager.java
@@ -41,13 +41,11 @@ import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
 import java.util.stream.Collectors;
 
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
+import lombok.extern.slf4j.Slf4j;
 
+@Slf4j
 public class ConsumerManager {
-
-    private static final Logger LOGGER = 
LoggerFactory.getLogger(ConsumerManager.class);
-
+    
     private final EventMeshGrpcServer eventMeshGrpcServer;
 
     private final ScheduledExecutorService scheduledExecutorService = 
Executors.newSingleThreadScheduledExecutor();
@@ -63,12 +61,12 @@ public class ConsumerManager {
     }
 
     public void init() throws Exception {
-        LOGGER.info("Grpc ConsumerManager initialized......");
+        log.info("Grpc ConsumerManager initialized......");
     }
 
     public void start() throws Exception {
         startClientCheck();
-        LOGGER.info("Grpc ConsumerManager started......");
+        log.info("Grpc ConsumerManager started......");
     }
 
     public void shutdown() throws Exception {
@@ -76,7 +74,7 @@ public class ConsumerManager {
             consumer.shutdown();
         }
         scheduledExecutorService.shutdown();
-        LOGGER.info("Grpc ConsumerManager shutdown......");
+        log.info("Grpc ConsumerManager shutdown......");
     }
 
     public EventMeshConsumer getEventMeshConsumer(String consumerGroup) {
@@ -163,7 +161,7 @@ public class ConsumerManager {
             }
         }
 
-        if (localClients.size() == 0) {
+        if (localClients.isEmpty()) {
             clientTable.remove(consumerGroup);
         }
     }
@@ -197,13 +195,13 @@ public class ConsumerManager {
         int clientTimeout = 
eventMeshGrpcServer.getEventMeshGrpcConfiguration().eventMeshSessionExpiredInMills;
         if (clientTimeout > 0) {
             scheduledExecutorService.scheduleAtFixedRate(() -> {
-                LOGGER.info("grpc client info check");
+                log.info("grpc client info check");
                 List<ConsumerGroupClient> clientList = new LinkedList<>();
                 for (List<ConsumerGroupClient> clients : clientTable.values()) 
{
                     clientList.addAll(clients);
                 }
-                if (LOGGER.isDebugEnabled()) {
-                    LOGGER.debug("total number of ConsumerGroupClients: {}", 
clientList.size());
+                if (log.isDebugEnabled()) {
+                    log.debug("total number of ConsumerGroupClients: {}", 
clientList.size());
                 }
 
                 if (clientList.isEmpty()) {
@@ -212,7 +210,7 @@ public class ConsumerManager {
                 Set<String> consumerGroupRestart = new HashSet<>();
                 for (ConsumerGroupClient client : clientList) {
                     if (System.currentTimeMillis() - 
client.getLastUpTime().getTime() > clientTimeout) {
-                        LOGGER.warn("client {} lastUpdate time {} over three 
heartbeat cycles. Removing it",
+                        log.warn("client {} lastUpdate time {} over three 
heartbeat cycles. Removing it",
                                 JsonUtils.serialize(client), 
client.getLastUpTime());
                         String consumerGroup = client.getConsumerGroup();
                         EventMeshConsumer consumer = 
getEventMeshConsumer(consumerGroup);
@@ -229,7 +227,7 @@ public class ConsumerManager {
                     try {
                         restartEventMeshConsumer(consumerGroup);
                     } catch (Exception e) {
-                        LOGGER.error("Error in restarting EventMeshConsumer 
[{}]", consumerGroup, e);
+                        log.error("Error in restarting EventMeshConsumer 
[{}]", consumerGroup, e);
                     }
                 }
             }, 10000, 10000, TimeUnit.MILLISECONDS);


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to