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 2ee71b521 fix issue2653
     new 374e7f870 Merge pull request #2683 from jonyangx/issue2653
2ee71b521 is described below

commit 2ee71b521a1b69e9601ddcd199afb3259a842d46
Author: jonyangx <[email protected]>
AuthorDate: Sat Dec 24 10:19:49 2022 +0800

    fix issue2653
---
 .../client/catalog/EventMeshCatalogClient.java     | 45 ++++++++++++----------
 1 file changed, 24 insertions(+), 21 deletions(-)

diff --git 
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/catalog/EventMeshCatalogClient.java
 
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/catalog/EventMeshCatalogClient.java
index 80f7e2017..13097716f 100644
--- 
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/catalog/EventMeshCatalogClient.java
+++ 
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/catalog/EventMeshCatalogClient.java
@@ -41,39 +41,42 @@ import io.grpc.ManagedChannel;
 import io.grpc.ManagedChannelBuilder;
 
 public class EventMeshCatalogClient {
-    private static final Logger logger = 
LoggerFactory.getLogger(EventMeshCatalogClient.class);
-    private final EventMeshCatalogClientConfig clientConfig;
-    private final EventMeshGrpcConsumer eventMeshGrpcConsumer;
-    private final List<SubscriptionItem> subscriptionItems = new ArrayList<>();
+    private static final Logger LOGGER = 
LoggerFactory.getLogger(EventMeshCatalogClient.class);
+    private final transient EventMeshCatalogClientConfig clientConfig;
+    private final transient EventMeshGrpcConsumer eventMeshGrpcConsumer;
+    private final transient List<SubscriptionItem> subscriptionItems = new 
ArrayList<>();
 
-    public EventMeshCatalogClient(EventMeshCatalogClientConfig clientConfig, 
EventMeshGrpcConsumer eventMeshGrpcConsumer) {
+    public EventMeshCatalogClient(EventMeshCatalogClientConfig clientConfig,
+                                  EventMeshGrpcConsumer eventMeshGrpcConsumer) 
{
         this.clientConfig = clientConfig;
         this.eventMeshGrpcConsumer = eventMeshGrpcConsumer;
     }
 
     public void init() throws Exception {
         Selector selector = 
SelectorFactory.get(clientConfig.getSelectorType());
-        AssertUtils.notNull(selector, String.format("selector=%s not 
register.please check it.", clientConfig.getSelectorType()));
+        AssertUtils.notNull(selector, String.format("selector=%s not 
register.please check it.",
+                clientConfig.getSelectorType()));
         ServiceInstance instance = 
selector.selectOne(clientConfig.getServerName());
         AssertUtils.notNull(instance, "catalog server is not running.please 
check it.");
+
         ManagedChannel channel = 
ManagedChannelBuilder.forAddress(instance.getHost(), instance.getPort())
-            .usePlaintext().build();
+                .usePlaintext().build();
         CatalogGrpc.CatalogBlockingStub catalogClient = 
CatalogGrpc.newBlockingStub(channel);
-        QueryOperationsRequest request = 
QueryOperationsRequest.newBuilder().setServiceName(clientConfig.getAppServerName()).build();
-        List<Operation> operations;
-        try {
-            QueryOperationsResponse response = 
catalogClient.queryOperations(request);
-            logger.info("received response: {}", response.toString());
-            operations = response.getOperationsList();
-            if (CollectionUtils.isEmpty(operations)) {
-                return;
-            }
-        } catch (Exception e) {
-            logger.error("queryOperations error {}", e.getMessage());
-            throw e;
+
+        QueryOperationsRequest request = QueryOperationsRequest.newBuilder()
+                .setServiceName(clientConfig.getAppServerName()).build();
+        QueryOperationsResponse response = 
catalogClient.queryOperations(request);
+        if (LOGGER.isInfoEnabled()) {
+            LOGGER.info("received response: {}", response.toString());
         }
+
+        List<Operation> operations = response.getOperationsList();
+        if (CollectionUtils.isEmpty(operations)) {
+            return;
+        }
+
         for (Operation operation : operations) {
-            if (!operation.getType().equals("subscribe")) {
+            if (!"subscribe".equals(operation.getType())) {
                 continue;
             }
             SubscriptionItem subscriptionItem = new SubscriptionItem();
@@ -86,7 +89,7 @@ public class EventMeshCatalogClient {
     }
 
     public void destroy() {
-        if (subscriptionItems.isEmpty()) {
+        if (subscriptionItems == null || subscriptionItems.isEmpty()) {
             return;
         }
         eventMeshGrpcConsumer.unsubscribe(subscriptionItems);


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

Reply via email to