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]