Fix relational tablet batch redirection and empty tablets
diff --git a/iotdb-client/session/src/main/java/org/apache/iotdb/session/Session.java b/iotdb-client/session/src/main/java/org/apache/iotdb/session/Session.java index 63e11b9..edb6690 100644 --- a/iotdb-client/session/src/main/java/org/apache/iotdb/session/Session.java +++ b/iotdb-client/session/src/main/java/org/apache/iotdb/session/Session.java
@@ -2874,6 +2874,19 @@ */ public void insertRelationalTablets(List<Tablet> tablets) throws IoTDBConnectionException, StatementExecutionException { + if (tablets.isEmpty()) { + throw new BatchExecutionException(SessionMessages.NO_TABLET_INSERTING); + } + final List<Tablet> nonEmptyTablets = new ArrayList<>(tablets.size()); + for (final Tablet tablet : tablets) { + if (tablet.getRowSize() > 0) { + nonEmptyTablets.add(tablet); + } + } + if (nonEmptyTablets.isEmpty()) { + return; + } + tablets = nonEmptyTablets; final boolean recordMergeTabletsCost = enableMergeTablets; long mergeTabletsCost = 0; if (enableMergeTablets) { @@ -2881,9 +2894,6 @@ tablets = mergeRelationalTablets(tablets); mergeTabletsCost = System.nanoTime() - mergeTabletsStartTime; } - if (tablets.isEmpty()) { - throw new BatchExecutionException(SessionMessages.NO_TABLET_INSERTING); - } final long insertTabletsStartTime = System.nanoTime(); if (enableRedirection) { insertRelationalTabletsWithLeaderCache(tablets); @@ -2893,10 +2903,7 @@ for (Tablet tablet : tablets) { request.addToColumnCategoriesList(toEnumOrdinalsAsBytes(tablet.getColumnTypes())); } - try { - getDefaultSessionConnection().insertTablets(request); - } catch (RedirectException ignored) { - } + getDefaultSessionConnection().insertTabletsWithoutRedirect(request); } if (recordMergeTabletsCost) { recordMergeTabletsCost(mergeTabletsCost, System.nanoTime() - insertTabletsStartTime);
diff --git a/iotdb-client/session/src/main/java/org/apache/iotdb/session/SessionConnection.java b/iotdb-client/session/src/main/java/org/apache/iotdb/session/SessionConnection.java index 4c30ba4..42ded5a 100644 --- a/iotdb-client/session/src/main/java/org/apache/iotdb/session/SessionConnection.java +++ b/iotdb-client/session/src/main/java/org/apache/iotdb/session/SessionConnection.java
@@ -876,6 +876,11 @@ () -> insertTabletsInternal(request), () -> devices); } + protected void insertTabletsWithoutRedirect(TSInsertTabletsReq request) + throws IoTDBConnectionException, StatementExecutionException { + callWithRetryAndVerify(() -> insertTabletsInternal(request)); + } + private TSStatus insertTabletsInternal(TSInsertTabletsReq request) throws TException { request.setSessionId(sessionId); return client.insertTablets(request);
diff --git a/iotdb-client/session/src/test/java/org/apache/iotdb/session/SessionConnectionTest.java b/iotdb-client/session/src/test/java/org/apache/iotdb/session/SessionConnectionTest.java index fea7846..5c3adb8 100644 --- a/iotdb-client/session/src/test/java/org/apache/iotdb/session/SessionConnectionTest.java +++ b/iotdb-client/session/src/test/java/org/apache/iotdb/session/SessionConnectionTest.java
@@ -352,6 +352,22 @@ } @Test + public void testInsertTabletsWithoutRedirectIgnoresRowRedirectStatus() + throws IoTDBConnectionException, StatementExecutionException, TException { + final TSStatus first = new TSStatus(TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode()); + first.setRedirectNode(new TEndPoint("127.0.0.2", 6667)); + final TSStatus second = new TSStatus(TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode()); + second.setRedirectNode(new TEndPoint("127.0.0.3", 6667)); + final TSStatus status = new TSStatus(TSStatusCode.REDIRECTION_RECOMMEND.getStatusCode()); + status.setSubStatus(Arrays.asList(first, second)); + Mockito.when(client.insertTablets(any())).thenReturn(status); + + sessionConnection.insertTabletsWithoutRedirect(new TSInsertTabletsReq()); + + Mockito.verify(client).insertTablets(any(TSInsertTabletsReq.class)); + } + + @Test public void testDeleteEmptyTimeseries() throws IoTDBConnectionException, StatementExecutionException, TException { sessionConnection.deleteTimeseries(Collections.emptyList());
diff --git a/iotdb-client/session/src/test/java/org/apache/iotdb/session/SessionTest.java b/iotdb-client/session/src/test/java/org/apache/iotdb/session/SessionTest.java index ac5959b..05d5f0a 100644 --- a/iotdb-client/session/src/test/java/org/apache/iotdb/session/SessionTest.java +++ b/iotdb-client/session/src/test/java/org/apache/iotdb/session/SessionTest.java
@@ -1369,6 +1369,57 @@ } @Test + public void testInsertRelationalTabletsWithRedirectionDisabledIgnoresRedirect() throws Exception { + final List<String> measurements = Arrays.asList("tag1", "s1"); + final Tablet tablet = + new Tablet( + "table1", + measurements, + Arrays.asList(TSDataType.STRING, TSDataType.INT64), + Arrays.asList(ColumnCategory.TAG, ColumnCategory.FIELD), + 2); + tablet.addTimestamp(0, 1L); + tablet.addValue("tag1", 0, "d1"); + tablet.addValue("s1", 0, 11L); + tablet.addTimestamp(1, 2L); + tablet.addValue("tag1", 1, "d2"); + tablet.addValue("s1", 1, 22L); + Whitebox.setInternalState(session, "enableRedirection", false); + + ((Session) session).insertRelationalTablets(Collections.singletonList(tablet)); + + Mockito.verify(sessionConnection).insertTabletsWithoutRedirect(any(TSInsertTabletsReq.class)); + Mockito.verify(sessionConnection, Mockito.never()).insertTablets(any(TSInsertTabletsReq.class)); + Mockito.verify(sessionConnection, Mockito.never()) + .insertTablets(any(TSInsertTabletsReq.class), anyList()); + } + + @Test + public void testInsertRelationalTabletsIgnoresAllEmptyTablets() throws Exception { + final Tablet first = + new Tablet( + "table1", + Arrays.asList("tag1", "s1"), + Arrays.asList(TSDataType.STRING, TSDataType.INT64), + Arrays.asList(ColumnCategory.TAG, ColumnCategory.FIELD), + 1); + final Tablet second = + new Tablet( + "table1", + Arrays.asList("tag1", "s1"), + Arrays.asList(TSDataType.STRING, TSDataType.INT64), + Arrays.asList(ColumnCategory.TAG, ColumnCategory.FIELD), + 1); + Whitebox.setInternalState(session, "enableRedirection", true); + + ((Session) session).insertRelationalTablets(Arrays.asList(first, second)); + + Mockito.verify(sessionConnection, Mockito.never()) + .insertTablets(any(TSInsertTabletsReq.class), anyList()); + Mockito.verify(sessionConnection, Mockito.never()).insertTablets(any(TSInsertTabletsReq.class)); + } + + @Test public void testSplitRelationalTabletAllocatesOncePerEndpoint() throws Exception { final int rowCount = 1000; final List<String> measurements = Arrays.asList("tag1", "s1");