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

haonan pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/master by this push:
     new 23bb94e2e5 [IOTDB-3271] Wrong multi IP for one sender in receiver 
(#6037)
23bb94e2e5 is described below

commit 23bb94e2e5874cf657bc3720b0e3b18cd2f64b78
Author: yschengzi <[email protected]>
AuthorDate: Mon May 30 12:45:57 2022 +0800

    [IOTDB-3271] Wrong multi IP for one sender in receiver (#6037)
---
 .../db/integration/sync/IoTDBSyncReceiverIT.java    |  5 ++---
 .../db/sync/sender/service/TransportHandler.java    |  5 ++---
 .../db/sync/transport/client/TransportClient.java   | 21 ++++++---------------
 .../db/sync/transport/TransportServiceTest.java     |  8 +++++---
 4 files changed, 15 insertions(+), 24 deletions(-)

diff --git 
a/integration/src/test/java/org/apache/iotdb/db/integration/sync/IoTDBSyncReceiverIT.java
 
b/integration/src/test/java/org/apache/iotdb/db/integration/sync/IoTDBSyncReceiverIT.java
index 97ee5e9a8b..7f296c5263 100644
--- 
a/integration/src/test/java/org/apache/iotdb/db/integration/sync/IoTDBSyncReceiverIT.java
+++ 
b/integration/src/test/java/org/apache/iotdb/db/integration/sync/IoTDBSyncReceiverIT.java
@@ -57,7 +57,6 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.io.File;
-import java.net.InetAddress;
 import java.net.Socket;
 import java.util.ArrayList;
 import java.util.Arrays;
@@ -108,8 +107,8 @@ public class IoTDBSyncReceiverIT {
       Assert.fail("Failed to start pipe server because " + e.getMessage());
     }
     Pipe pipe = new TsFilePipe(createdTime1, pipeName1, null, 0, false);
-    client = new TransportClient(pipe, "127.0.0.1", 6670);
-    remoteIp1 = InetAddress.getLocalHost().getHostAddress();
+    client = new TransportClient(pipe, "127.0.0.1", 6670, "127.0.0.1");
+    remoteIp1 = "127.0.0.1";
     client.handshake();
   }
 
diff --git 
a/server/src/main/java/org/apache/iotdb/db/sync/sender/service/TransportHandler.java
 
b/server/src/main/java/org/apache/iotdb/db/sync/sender/service/TransportHandler.java
index b6e9707de0..af49ab9d2e 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/sync/sender/service/TransportHandler.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/sync/sender/service/TransportHandler.java
@@ -62,7 +62,8 @@ public class TransportHandler {
   public TransportHandler(Pipe pipe, IoTDBPipeSink pipeSink) {
     this.pipeName = pipe.getName();
     this.createTime = pipe.getCreateTime();
-    this.transportClient = new TransportClient(pipe, pipeSink.getIp(), 
pipeSink.getPort());
+    this.localIP = getLocalIP(pipeSink);
+    this.transportClient = new TransportClient(pipe, pipeSink.getIp(), 
pipeSink.getPort(), localIP);
 
     this.transportExecutorService =
         IoTDBThreadPoolFactory.newSingleThreadExecutor(
@@ -70,8 +71,6 @@ public class TransportHandler {
     this.heartbeatExecutorService =
         IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor(
             ThreadName.SYNC_SENDER_HEARTBEAT.getName() + "-" + pipeName);
-
-    this.localIP = getLocalIP(pipeSink);
   }
 
   private String getLocalIP(IoTDBPipeSink pipeSink) {
diff --git 
a/server/src/main/java/org/apache/iotdb/db/sync/transport/client/TransportClient.java
 
b/server/src/main/java/org/apache/iotdb/db/sync/transport/client/TransportClient.java
index c495a7b256..7613f3f8ec 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/sync/transport/client/TransportClient.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/sync/transport/client/TransportClient.java
@@ -56,8 +56,6 @@ import java.io.FileReader;
 import java.io.IOException;
 import java.io.InputStream;
 import java.io.RandomAccessFile;
-import java.net.InetAddress;
-import java.net.UnknownHostException;
 import java.nio.ByteBuffer;
 import java.security.MessageDigest;
 import java.security.NoSuchAlgorithmException;
@@ -80,6 +78,7 @@ public class TransportClient implements ITransportClient {
   private TransportService.Client serviceClient = null;
 
   private String ipAddress;
+  private String localIP;
 
   private int port;
 
@@ -87,11 +86,12 @@ public class TransportClient implements ITransportClient {
 
   private Pipe pipe;
 
-  public TransportClient(Pipe pipe, String ipAddress, int port) {
+  public TransportClient(Pipe pipe, String ipAddress, int port, String 
localIP) {
     RpcTransportFactory.setThriftMaxFrameSize(config.getThriftMaxFrameSize());
 
     this.pipe = pipe;
     this.ipAddress = ipAddress;
+    this.localIP = localIP;
     this.port = port;
   }
 
@@ -148,10 +148,7 @@ public class TransportClient implements ITransportClient {
 
       identityInfo =
           new IdentityInfo(
-              InetAddress.getLocalHost().getHostAddress(),
-              pipe.getName(),
-              pipe.getCreateTime(),
-              config.getIoTDBMajorVersion());
+              localIP, pipe.getName(), pipe.getCreateTime(), 
config.getIoTDBMajorVersion());
       TransportStatus status = serviceClient.handshake(identityInfo);
       if (status.code != SUCCESS_CODE) {
         throw new SyncConnectionException(
@@ -160,9 +157,6 @@ public class TransportClient implements ITransportClient {
     } catch (TException e) {
       logger.warn("Cannot connect to the receiver. ", e);
       return false;
-    } catch (UnknownHostException e) {
-      logger.warn("Cannot confirm identity with the receiver. ", e);
-      throw new SyncConnectionException(String.format("Get local host error, 
because %s.", e), e);
     }
     return true;
   }
@@ -460,10 +454,7 @@ public class TransportClient implements ITransportClient {
           .receiveMsg(
               heartbeat(
                   new SyncRequest(
-                      RequestType.START,
-                      pipe.getName(),
-                      InetAddress.getLocalHost().getHostAddress(),
-                      pipe.getCreateTime())));
+                      RequestType.START, pipe.getName(), localIP, 
pipe.getCreateTime())));
       while (!Thread.currentThread().isInterrupted()) {
         PipeData pipeData = pipe.take();
         if (!senderTransport(pipeData)) {
@@ -481,7 +472,7 @@ public class TransportClient implements ITransportClient {
       }
     } catch (InterruptedException e) {
       logger.info("Interrupted by pipe, exit transport.");
-    } catch (SyncConnectionException | UnknownHostException e) {
+    } catch (SyncConnectionException e) {
       logger.error(
           String.format("Connect to receiver %s:%d error, because %s.", 
ipAddress, port, e));
       SenderService.getInstance()
diff --git 
a/server/src/test/java/org/apache/iotdb/db/sync/transport/TransportServiceTest.java
 
b/server/src/test/java/org/apache/iotdb/db/sync/transport/TransportServiceTest.java
index f4e1742b7b..d4ea1d3ab4 100644
--- 
a/server/src/test/java/org/apache/iotdb/db/sync/transport/TransportServiceTest.java
+++ 
b/server/src/test/java/org/apache/iotdb/db/sync/transport/TransportServiceTest.java
@@ -52,7 +52,6 @@ import java.io.BufferedInputStream;
 import java.io.File;
 import java.io.FileInputStream;
 import java.io.FileWriter;
-import java.net.InetAddress;
 import java.security.MessageDigest;
 import java.util.ArrayList;
 import java.util.List;
@@ -71,7 +70,7 @@ public class TransportServiceTest {
 
   @Before
   public void setUp() throws Exception {
-    remoteIp1 = InetAddress.getLocalHost().getHostAddress();
+    remoteIp1 = "127.0.0.1";
     fileDir = new File(SyncPathUtil.getReceiverFileDataDir(pipeName1, 
remoteIp1, createdTime1));
     pipeDataQueue =
         PipeDataQueueFactory.getBufferedPipeDataQueue(
@@ -131,7 +130,10 @@ public class TransportServiceTest {
     Pipe pipe = new TsFilePipe(createdTime1, pipeName1, null, 0, false);
     TransportClient client =
         new TransportClient(
-            pipe, "127.0.0.1", 
IoTDBDescriptor.getInstance().getConfig().getPipeServerPort());
+            pipe,
+            "127.0.0.1",
+            IoTDBDescriptor.getInstance().getConfig().getPipeServerPort(),
+            "127.0.0.1");
     client.handshake();
     for (PipeData pipeData : pipeDataList) {
       client.senderTransport(pipeData);

Reply via email to