test(load): cover Sonar warning fixes
diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/load/LoadSingleTsFileNode.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/load/LoadSingleTsFileNode.java index ae4b753..79fd5e8 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/load/LoadSingleTsFileNode.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/load/LoadSingleTsFileNode.java
@@ -50,9 +50,7 @@ import java.util.Collections; import java.util.HashSet; import java.util.List; -import java.util.NoSuchElementException; import java.util.Objects; -import java.util.Optional; import java.util.Set; import java.util.function.Function; @@ -102,16 +100,10 @@ new ArrayList<>(resource.getDevices().size() << 1); for (final IDeviceID device : resource.getDevices()) { // iterating the index, must present - final Optional<Long> startTime = resource.getStartTime(device); - if (!startTime.isPresent()) { - throw new NoSuchElementException("No value present"); - } - final Optional<Long> endTime = resource.getEndTime(device); - if (!endTime.isPresent()) { - throw new NoSuchElementException("No value present"); - } - final TTimePartitionSlot startSlot = TimePartitionUtils.getTimePartitionSlot(startTime.get()); - final TTimePartitionSlot endSlot = TimePartitionUtils.getTimePartitionSlot(endTime.get()); + final long startTime = resource.getStartTime(device).orElseThrow(); + final long endTime = resource.getEndTime(device).orElseThrow(); + final TTimePartitionSlot startSlot = TimePartitionUtils.getTimePartitionSlot(startTime); + final TTimePartitionSlot endSlot = TimePartitionUtils.getTimePartitionSlot(endTime); slotList.add(new Pair<>(device, startSlot)); if (!startSlot.equals(endSlot)) { slotList.add(new Pair<>(device, endSlot));
diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java index c05aaab..69ae38d 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoader.java
@@ -203,7 +203,7 @@ if (!loadEntry.isPresent()) { return; } - final ActiveLoadPendingQueue.ActiveLoadEntry activeLoadEntry = loadEntry.get(); + final ActiveLoadPendingQueue.ActiveLoadEntry activeLoadEntry = loadEntry.orElseThrow(); try { final TSStatus result = loadTsFile(activeLoadEntry, session);
diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/node/load/LoadTsFileNodeTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/node/load/LoadTsFileNodeTest.java index 37cedc5..757a6a3 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/node/load/LoadTsFileNodeTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/planner/node/load/LoadTsFileNodeTest.java
@@ -25,9 +25,11 @@ import org.apache.iotdb.db.queryengine.plan.planner.plan.node.load.LoadTsFilePieceNode; import org.apache.iotdb.db.storageengine.dataregion.modification.ModificationFile; import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource; +import org.apache.iotdb.db.storageengine.dataregion.tsfile.timeindex.ArrayDeviceTimeIndex; import org.apache.iotdb.db.storageengine.load.util.LoadUtil; import org.apache.tsfile.exception.NotImplementedException; +import org.apache.tsfile.file.metadata.IDeviceID; import org.junit.Assert; import org.junit.Test; @@ -94,6 +96,20 @@ } @Test + public void testNeedDecodeTsFileReadsDeviceTimeIndex() { + final TsFileResource resource = new TsFileResource(new File("load.tsfile")); + resource.setTimeIndex(new ArrayDeviceTimeIndex()); + final IDeviceID device = IDeviceID.Factory.DEFAULT_FACTORY.create("root.sg.d1"); + resource.updateStartTime(device, 0L); + resource.updateEndTime(device, 1L); + + final LoadSingleTsFileNode node = + new LoadSingleTsFileNode(new PlanNodeId(""), resource, false, null, false, 0L, false); + + Assert.assertFalse(node.needDecodeTsFile(slotList -> Collections.emptyList())); + } + + @Test public void testCleanContinuesAfterOneFileCannotBeDeleted() throws Exception { final File tempDir = Files.createTempDirectory("load-node-clean").toFile(); try {
diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileSchedulerTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileSchedulerTest.java index 1b738f6..ee64e80 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileSchedulerTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileSchedulerTest.java
@@ -19,7 +19,14 @@ package org.apache.iotdb.db.queryengine.plan.scheduler.load; +import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId; +import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType; +import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet; +import org.apache.iotdb.common.rpc.thrift.TTimePartitionSlot; import org.apache.iotdb.commons.client.IClientManager; +import org.apache.iotdb.commons.partition.DataPartition; +import org.apache.iotdb.commons.queryengine.common.SessionInfo; +import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNodeId; import org.apache.iotdb.db.queryengine.common.MPPQueryContext; import org.apache.iotdb.db.queryengine.common.PlanFragmentId; import org.apache.iotdb.db.queryengine.execution.QueryStateMachine; @@ -29,8 +36,11 @@ import org.apache.iotdb.db.queryengine.plan.planner.plan.SubPlan; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.load.LoadSingleTsFileNode; import org.apache.iotdb.db.queryengine.plan.statement.crud.LoadTsFileStatement; +import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource; import org.apache.iotdb.db.storageengine.load.memory.LoadTsFileDataCacheMemoryBlock; +import org.apache.iotdb.db.storageengine.load.splitter.ChunkData; +import org.apache.tsfile.file.metadata.IDeviceID; import org.junit.Assert; import org.junit.Before; import org.junit.Test; @@ -41,9 +51,15 @@ import java.lang.reflect.Constructor; import java.lang.reflect.Field; import java.lang.reflect.Method; +import java.util.Collections; +import java.util.List; +import static org.mockito.ArgumentMatchers.anyList; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; public class LoadTsFileSchedulerTest { @@ -58,6 +74,7 @@ when(distributedQueryPlan.getRootSubPlan()).thenReturn(subPlan); when(subPlan.getPlanFragment()).thenReturn(planFragment); when(planFragment.getId()).thenReturn(new PlanFragmentId("test", 0)); + when(distributedQueryPlan.getInstances()).thenReturn(Collections.emptyList()); } @Test @@ -160,4 +177,73 @@ Assert.assertEquals(0L, getMemoryUsageMethod.invoke(memoryBlock)); Assert.assertEquals(0L, dataSizeField.getLong(dataManager)); } + + @Test + public void testRouteChunkDataDeduplicatesPartitionSlots() throws Exception { + final IPartitionFetcher partitionFetcher = mock(IPartitionFetcher.class); + final DataPartition dataPartition = mock(DataPartition.class); + final MPPQueryContext queryContext = mock(MPPQueryContext.class); + final SessionInfo sessionInfo = mock(SessionInfo.class); + when(queryContext.getSession()).thenReturn(sessionInfo); + when(sessionInfo.getUserName()).thenReturn("root"); + when(partitionFetcher.getOrCreateDataPartition(anyList(), eq("root"))) + .thenReturn(dataPartition); + + final IDeviceID device = IDeviceID.Factory.DEFAULT_FACTORY.create("root.sg.d1"); + final TTimePartitionSlot timePartitionSlot = new TTimePartitionSlot(0L); + final TRegionReplicaSet replicaSet = + new TRegionReplicaSet( + new TConsensusGroupId(TConsensusGroupType.DataRegion, 0), Collections.emptyList()); + when(dataPartition.getDataRegionReplicaSetForWriting(device, timePartitionSlot)) + .thenReturn(replicaSet); + + final LoadTsFileScheduler scheduler = + new LoadTsFileScheduler( + distributedQueryPlan, + queryContext, + mock(QueryStateMachine.class), + mock(IClientManager.class), + partitionFetcher, + false); + final LoadSingleTsFileNode node = mock(LoadSingleTsFileNode.class); + final TsFileResource resource = mock(TsFileResource.class); + when(node.getPlanNodeId()).thenReturn(new PlanNodeId("test")); + when(node.getTsFileResource()).thenReturn(resource); + when(resource.getTsFile()).thenReturn(new File("test.tsfile")); + + final Class<?> dataManagerClass = + Class.forName(LoadTsFileScheduler.class.getName() + "$TsFileDataManager"); + final Constructor<?> dataManagerConstructor = + dataManagerClass.getDeclaredConstructor( + LoadTsFileScheduler.class, + LoadSingleTsFileNode.class, + LoadTsFileDataCacheMemoryBlock.class); + dataManagerConstructor.setAccessible(true); + final Object dataManager = + dataManagerConstructor.newInstance( + scheduler, node, mock(LoadTsFileDataCacheMemoryBlock.class)); + + final ChunkData firstChunk = mock(ChunkData.class); + when(firstChunk.getDevice()).thenReturn(device); + when(firstChunk.getTimePartitionSlot()).thenReturn(timePartitionSlot); + when(firstChunk.getDataSize()).thenReturn(1L); + final ChunkData secondChunk = mock(ChunkData.class); + when(secondChunk.getDevice()).thenReturn(device); + when(secondChunk.getTimePartitionSlot()).thenReturn(timePartitionSlot); + when(secondChunk.getDataSize()).thenReturn(1L); + + final Field chunkDataField = dataManagerClass.getDeclaredField("nonDirectionalChunkData"); + chunkDataField.setAccessible(true); + @SuppressWarnings("unchecked") + final List<ChunkData> chunkData = (List<ChunkData>) chunkDataField.get(dataManager); + chunkData.add(firstChunk); + chunkData.add(secondChunk); + + final Method routeChunkData = dataManagerClass.getDeclaredMethod("routeChunkData"); + routeChunkData.setAccessible(true); + routeChunkData.invoke(dataManager); + + verify(dataPartition, times(1)).getDataRegionReplicaSetForWriting(device, timePartitionSlot); + Assert.assertTrue(chunkData.isEmpty()); + } }
diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoaderTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoaderTest.java index aa81de4..d1ffe6d 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoaderTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/active/ActiveLoadTsFileLoaderTest.java
@@ -38,6 +38,8 @@ import java.lang.reflect.Field; import java.lang.reflect.Method; import java.nio.file.Files; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; public class ActiveLoadTsFileLoaderTest { @@ -45,13 +47,16 @@ private File tempDir; private String originalFailDir; private NodeStatus originalNodeStatus; + private int originalDataNodeId; @Before public void setUp() throws Exception { tempDir = Files.createTempDirectory("active-load-retry").toFile(); originalFailDir = config.getLoadActiveListeningFailDir(); originalNodeStatus = CommonDescriptor.getInstance().getConfig().getNodeStatus(); + originalDataNodeId = config.getDataNodeId(); CommonDescriptor.getInstance().getConfig().setNodeStatus(NodeStatus.Running); + config.setDataNodeId(0); config.setLoadActiveListeningFailDir(new File(tempDir, "failed").getAbsolutePath()); } @@ -59,6 +64,7 @@ public void tearDown() { config.setLoadActiveListeningFailDir(originalFailDir); CommonDescriptor.getInstance().getConfig().setNodeStatus(originalNodeStatus); + config.setDataNodeId(originalDataNodeId); deleteRecursively(tempDir); } @@ -127,6 +133,47 @@ Assert.assertTrue(pendingQueue.isEmpty()); } + @Test + public void testMissingPendingFileIsRemovedFromLoading() throws Exception { + final ActiveLoadTsFileLoader loader = new ActiveLoadTsFileLoader(); + final Field pendingQueueField = ActiveLoadTsFileLoader.class.getDeclaredField("pendingQueue"); + pendingQueueField.setAccessible(true); + final ActiveLoadPendingQueue pendingQueue = + (ActiveLoadPendingQueue) pendingQueueField.get(loader); + final File missingTsFile = new File(tempDir, "missing.tsfile"); + Assert.assertTrue( + pendingQueue.enqueue( + missingTsFile.getAbsolutePath(), tempDir.getAbsolutePath(), false, false)); + + final AtomicReference<Throwable> failure = new AtomicReference<>(); + final Thread loadThread = + new Thread( + () -> { + try { + invokeTryLoadPendingTsFiles(loader); + } catch (final Throwable e) { + Throwable cause = e; + while (cause.getCause() != null) { + cause = cause.getCause(); + } + failure.set(cause); + } + }); + loadThread.start(); + + final long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5); + while (loader.isFilePendingOrLoading(missingTsFile) && System.nanoTime() < deadline) { + Thread.yield(); + } + loadThread.interrupt(); + loadThread.join(TimeUnit.SECONDS.toMillis(5)); + + Assert.assertFalse(loadThread.isAlive()); + Assert.assertNull(failure.get()); + Assert.assertFalse(loader.isFilePendingOrLoading(missingTsFile)); + Assert.assertTrue(pendingQueue.isEmpty()); + } + private File createTsFileWithCompanionFiles(final String fileName) throws Exception { final File tsFile = new File(tempDir, fileName); Assert.assertTrue(tsFile.createNewFile()); @@ -157,6 +204,12 @@ method.invoke(loader, entry, status); } + private void invokeTryLoadPendingTsFiles(final ActiveLoadTsFileLoader loader) throws Exception { + final Method method = ActiveLoadTsFileLoader.class.getDeclaredMethod("tryLoadPendingTsFiles"); + method.setAccessible(true); + method.invoke(loader); + } + private static void deleteRecursively(final File file) { if (file == null || !file.exists()) { return;