luchunliang commented on a change in pull request #2670:
URL: https://github.com/apache/incubator-inlong/pull/2670#discussion_r813489361



##########
File path: 
inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/DataProxyClusterServiceImpl.java
##########
@@ -264,16 +279,17 @@ public ProxyPulsarDTO getConfigV2(String 
dataproxyClusterName) {
          */

Review comment:
       It is not suggested to use "/* */ in code block.

##########
File path: 
inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/PulsarSink.java
##########
@@ -283,21 +291,28 @@ public void diffSetPublish(PulsarClientService 
pulsarClientService, Set<String>
      * @param originalCluster
      * @param endCluster
      */
-    public void diffRestartPulsarClient(Set<String> originalCluster, 
Set<String> endCluster) {
-        if (!originalCluster.equals(endCluster)) {
-            logger.info("pulsarConfig has changed, close current 
pulsarClientService and restart");
-            pulsarClientService.close();
-
-            pulsarCluster = configManager.getPulsarUrl2Token();
-            configManager.getPulsarConfig().setUrl2token(pulsarCluster);
-            pulsarClientService.initCreateConnection(this);
-            try {
-                initTopicSet(pulsarClientService, new 
HashSet<String>(topicProperties.values()));
-            } catch (Exception e) {
-                logger.info("pulsar sink restart, publish topic fail.", e);
-            }
+    public void diffUpdatePulsarClient(PulsarClientService 
pulsarClientService, Map<String, String> originalCluster,
+                                       Map<String, String> endCluster) {
+        MapDifference<String, String> mapDifference = 
Maps.difference(originalCluster, endCluster);
+        if (mapDifference.areEqual()) {
+            return;
+        }
 
+        logger.info("pulsarConfig has changed, close unused url clients and 
start new url clients");
+        Map<String, String> needToStart = new HashMap<>();
+        Map<String, String> needToClose = new HashMap<>();
+        needToClose.putAll(mapDifference.entriesOnlyOnLeft());
+        needToStart.putAll(mapDifference.entriesOnlyOnRight());
+        Map<String, MapDifference.ValueDifference<String>> differentToken = 
mapDifference.entriesDiffering();
+        for (String url : differentToken.keySet()) {
+            needToClose.put(url, originalCluster.get(url));
+            needToStart.put(url, endCluster.get(url));//token changed
         }
+
+        pulsarClientService.updatePulsarClients(this, needToClose, needToStart,

Review comment:
       It is not enough that pulsar client only use two status "start and 
close".
   Closing status is necessary for waiting the completion of sending messages.
   

##########
File path: 
inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/pulsar/PulsarClientService.java
##########
@@ -299,6 +300,65 @@ private void destroyConnection() {
         logger.debug("closed meta producer");
     }
 
+    /**
+     * close pulsarClients(the related url is removed); start pulsarClients 
for new url, and create producers for them
+     *
+     * @param callBack
+     * @param needToClose url-token map
+     * @param needToStart url-token map
+     * @param topicSet    for new pulsarClient, create these topics' producers
+     */
+    public void updatePulsarClients(CreatePulsarClientCallBack callBack, 
Map<String, String> needToClose,
+                                    Map<String, String> needToStart, 
Set<String> topicSet) {
+        // close
+        for (String url : needToClose.keySet()) {
+            PulsarClient pulsarClient = pulsarClients.get(url);
+            if (pulsarClient != null) {
+                try {
+                    pulsarClient.shutdown();

Review comment:
       Miss the close operation of topic producer, it will lost the messages in 
the producer buffer.

##########
File path: 
inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/pulsar/PulsarClientService.java
##########
@@ -299,6 +300,65 @@ private void destroyConnection() {
         logger.debug("closed meta producer");
     }
 
+    /**
+     * close pulsarClients(the related url is removed); start pulsarClients 
for new url, and create producers for them
+     *
+     * @param callBack
+     * @param needToClose url-token map
+     * @param needToStart url-token map
+     * @param topicSet    for new pulsarClient, create these topics' producers
+     */
+    public void updatePulsarClients(CreatePulsarClientCallBack callBack, 
Map<String, String> needToClose,
+                                    Map<String, String> needToStart, 
Set<String> topicSet) {
+        // close
+        for (String url : needToClose.keySet()) {
+            PulsarClient pulsarClient = pulsarClients.get(url);
+            if (pulsarClient != null) {
+                try {
+                    pulsarClient.shutdown();
+                    pulsarClients.remove(url);
+                } catch (PulsarClientException e) {
+                    logger.error("shutdown pulsarClient error in PulsarSink, 
PulsarClientException {}",
+                            e.getMessage());
+                } catch (Exception e) {
+                    logger.error("shutdown pulsarClient error in PulsarSink, 
ex {}", e.getMessage());
+                }
+            }
+        }
+        // new pulsarClient
+        for (Map.Entry<String, String> entry : needToStart.entrySet()) {
+            String url = entry.getKey();
+            String token = entry.getValue();
+            try {
+                if (logger.isDebugEnabled()) {
+                    logger.debug("url = {}, token = {}", url, token);
+                }
+                PulsarClient client = initPulsarClient(url, token);
+                pulsarClients.put(url, client);
+                callBack.handleCreateClientSuccess(url);

Review comment:
       The operation of "pulsarClients.put(url, client)" must move after the 
init opeartion of topic producer.
   It maybe have crash of the init opeartion of topic producer between 
sendMessage and updatePulsarClients.

##########
File path: 
inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/pulsar/PulsarClientService.java
##########
@@ -299,6 +300,65 @@ private void destroyConnection() {
         logger.debug("closed meta producer");
     }
 
+    /**
+     * close pulsarClients(the related url is removed); start pulsarClients 
for new url, and create producers for them

Review comment:
       This comment is not fit for the method name "updatePulsarClients"




-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


Reply via email to