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 60242a4ec update log
     new 650de8c8b Merge pull request #2744 from weihubeats/ConsumerService
60242a4ec is described below

commit 60242a4ec07cc1565951a147eadade73a29c5ff6
Author: weihu <[email protected]>
AuthorDate: Fri Dec 30 14:37:11 2022 +0800

    update log
---
 .../protocol/grpc/service/ConsumerService.java     | 28 ++++++++++------------
 1 file changed, 13 insertions(+), 15 deletions(-)

diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/service/ConsumerService.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/service/ConsumerService.java
index 280a2d597..ec77420b0 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/service/ConsumerService.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/service/ConsumerService.java
@@ -31,14 +31,12 @@ import 
org.apache.eventmesh.runtime.core.protocol.grpc.processor.UnsubscribeProc
 
 import java.util.concurrent.ThreadPoolExecutor;
 
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
 import io.grpc.stub.StreamObserver;
 
-public class ConsumerService extends 
ConsumerServiceGrpc.ConsumerServiceImplBase {
+import lombok.extern.slf4j.Slf4j;
 
-    private final Logger logger = 
LoggerFactory.getLogger(ConsumerService.class);
+@Slf4j
+public class ConsumerService extends 
ConsumerServiceGrpc.ConsumerServiceImplBase {
 
     private final EventMeshGrpcServer eventMeshGrpcServer;
 
@@ -55,7 +53,7 @@ public class ConsumerService extends 
ConsumerServiceGrpc.ConsumerServiceImplBase
     }
 
     public void subscribe(Subscription request, StreamObserver<Response> 
responseObserver) {
-        logger.info("cmd={}|{}|client2eventMesh|from={}|to={}",
+        log.info("cmd={}|{}|client2eventMesh|from={}|to={}",
             "subscribe", EventMeshConstants.PROTOCOL_GRPC,
             request.getHeader().getIp(), 
eventMeshGrpcServer.getEventMeshGrpcConfiguration().eventMeshIp);
         eventMeshGrpcServer.getMetricsMonitor().recordReceiveMsgFromClient();
@@ -66,7 +64,7 @@ public class ConsumerService extends 
ConsumerServiceGrpc.ConsumerServiceImplBase
             try {
                 subscribeProcessor.process(request, emitter);
             } catch (Exception e) {
-                logger.error("Error code {}, error message {}", 
StatusCode.EVENTMESH_SUBSCRIBE_ERR.getRetCode(),
+                log.error("Error code {}, error message {}", 
StatusCode.EVENTMESH_SUBSCRIBE_ERR.getRetCode(),
                     StatusCode.EVENTMESH_SUBSCRIBE_ERR.getErrMsg(), e);
                 
ServiceUtils.sendRespAndDone(StatusCode.EVENTMESH_SUBSCRIBE_ERR, 
e.getMessage(), emitter);
             }
@@ -80,14 +78,14 @@ public class ConsumerService extends 
ConsumerServiceGrpc.ConsumerServiceImplBase
             @Override
             public void onNext(Subscription subscription) {
                 if (!subscription.getSubscriptionItemsList().isEmpty()) {
-                    logger.info("cmd={}|{}|client2eventMesh|from={}|to={}",
+                    log.info("cmd={}|{}|client2eventMesh|from={}|to={}",
                         "subscribeStream", EventMeshConstants.PROTOCOL_GRPC,
                         subscription.getHeader().getIp(), 
eventMeshGrpcServer.getEventMeshGrpcConfiguration().eventMeshIp);
                     
eventMeshGrpcServer.getMetricsMonitor().recordReceiveMsgFromClient();
 
                     handleSubscriptionStream(subscription, emitter);
                 } else {
-                    logger.info("cmd={}|{}|client2eventMesh|from={}|to={}",
+                    log.info("cmd={}|{}|client2eventMesh|from={}|to={}",
                         "reply-to-server", EventMeshConstants.PROTOCOL_GRPC,
                         subscription.getHeader().getIp(), 
eventMeshGrpcServer.getEventMeshGrpcConfiguration().eventMeshIp);
 
@@ -97,13 +95,13 @@ public class ConsumerService extends 
ConsumerServiceGrpc.ConsumerServiceImplBase
 
             @Override
             public void onError(Throwable t) {
-                logger.error("Receive error from client: " + t.getMessage());
+                log.error("Receive error from client: {}", t.getMessage());
                 emitter.onCompleted();
             }
 
             @Override
             public void onCompleted() {
-                logger.info("Client finish sending messages");
+                log.info("Client finish sending messages");
                 emitter.onCompleted();
             }
         };
@@ -115,7 +113,7 @@ public class ConsumerService extends 
ConsumerServiceGrpc.ConsumerServiceImplBase
             try {
                 streamProcessor.process(request, emitter);
             } catch (Exception e) {
-                logger.error("Error code {}, error message {}", 
StatusCode.EVENTMESH_SUBSCRIBE_ERR, e.getMessage(), e);
+                log.error("Error code {}, error message {}", 
StatusCode.EVENTMESH_SUBSCRIBE_ERR, e.getMessage(), e);
                 ServiceUtils.sendStreamRespAndDone(request.getHeader(), 
StatusCode.EVENTMESH_SUBSCRIBE_ERR, e.getMessage(), emitter);
             }
         });
@@ -127,7 +125,7 @@ public class ConsumerService extends 
ConsumerServiceGrpc.ConsumerServiceImplBase
             try {
                 
replyMessageProcessor.process(buildSimpleMessage(subscription), emitter);
             } catch (Exception e) {
-                logger.error("Error code {}, error message {}", 
StatusCode.EVENTMESH_SUBSCRIBE_ERR, e.getMessage(), e);
+                log.error("Error code {}, error message {}", 
StatusCode.EVENTMESH_SUBSCRIBE_ERR, e.getMessage(), e);
                 ServiceUtils.sendStreamRespAndDone(subscription.getHeader(), 
StatusCode.EVENTMESH_SUBSCRIBE_ERR, e.getMessage(), emitter);
             }
         });
@@ -148,7 +146,7 @@ public class ConsumerService extends 
ConsumerServiceGrpc.ConsumerServiceImplBase
     }
 
     public void unsubscribe(Subscription request, StreamObserver<Response> 
responseObserver) {
-        logger.info("cmd={}|{}|client2eventMesh|from={}|to={}",
+        log.info("cmd={}|{}|client2eventMesh|from={}|to={}",
             "unsubscribe", EventMeshConstants.PROTOCOL_GRPC,
             request.getHeader().getIp(), 
eventMeshGrpcServer.getEventMeshGrpcConfiguration().eventMeshIp);
         eventMeshGrpcServer.getMetricsMonitor().recordReceiveMsgFromClient();
@@ -159,7 +157,7 @@ public class ConsumerService extends 
ConsumerServiceGrpc.ConsumerServiceImplBase
             try {
                 unsubscribeProcessor.process(request, emitter);
             } catch (Exception e) {
-                logger.error("Error code {}, error message {}", 
StatusCode.EVENTMESH_UNSUBSCRIBE_ERR.getRetCode(),
+                log.error("Error code {}, error message {}", 
StatusCode.EVENTMESH_UNSUBSCRIBE_ERR.getRetCode(),
                     StatusCode.EVENTMESH_UNSUBSCRIBE_ERR.getErrMsg(), e);
                 
ServiceUtils.sendRespAndDone(StatusCode.EVENTMESH_UNSUBSCRIBE_ERR, 
e.getMessage(), emitter);
             }


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

Reply via email to