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);