This is an automated email from the ASF dual-hosted git repository. jt2594838 pushed a commit to branch fix_iot_memory_stuck in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 0fb9f773b531796f821a5fadb5f3d206d72a5bf0 Author: Tian Jiang <[email protected]> AuthorDate: Wed Aug 12 16:40:09 2026 +0800 fix: retry SyncStatus batch memory reservation --- .../consensus/iot/logdispatcher/SyncStatus.java | 13 +++-- .../iot/logdispatcher/SyncStatusTest.java | 64 ++++++++++++++++++++++ 2 files changed, 73 insertions(+), 4 deletions(-) diff --git a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/SyncStatus.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/SyncStatus.java index 08296aef1f5..1749384f549 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/SyncStatus.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/SyncStatus.java @@ -49,10 +49,15 @@ public class SyncStatus { * @throws InterruptedException */ public synchronized void addNextBatch(Batch batch) throws InterruptedException { - while ((pendingBatches.size() >= config.getReplication().getMaxPendingBatchesNum() - || !iotConsensusMemoryManager.reserve(batch)) - && !Thread.interrupted()) { - wait(); + while (true) { + while (pendingBatches.size() >= config.getReplication().getMaxPendingBatchesNum()) { + wait(); + } + if (iotConsensusMemoryManager.reserve(batch)) { + break; + } + // Memory may be freed by another SyncStatus, which cannot notify this monitor. + wait(Math.max(1, config.getReplication().getBasicRetryWaitTimeMs())); } if (LOGGER.isDebugEnabled()) { LOGGER.debug( diff --git a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/SyncStatusTest.java b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/SyncStatusTest.java index be81c69f7f7..c70b97e7131 100644 --- a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/SyncStatusTest.java +++ b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/SyncStatusTest.java @@ -21,6 +21,8 @@ package org.apache.iotdb.consensus.iot.logdispatcher; import org.apache.iotdb.common.rpc.thrift.TEndPoint; import org.apache.iotdb.commons.consensus.DataRegionId; +import org.apache.iotdb.commons.memory.AtomicLongMemoryBlock; +import org.apache.iotdb.commons.memory.IMemoryBlock; import org.apache.iotdb.consensus.common.Peer; import org.apache.iotdb.consensus.config.IoTConsensusConfig; import org.apache.iotdb.consensus.iot.thrift.TLogEntry; @@ -36,7 +38,14 @@ import java.io.IOException; import java.util.ArrayList; import java.util.List; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicBoolean; public class SyncStatusTest { @@ -242,4 +251,59 @@ public class SyncStatusTest { Assert.assertEquals( config.getReplication().getMaxPendingBatchesNum() + 1, status.getNextSendingIndex()); } + + @Test + public void testFirstBatchRetriesMemoryReservation() + throws InterruptedException, ExecutionException, TimeoutException { + IndexController controller = + new IndexController(storageDir.getAbsolutePath(), peer, 0, CHECK_POINT_GAP); + IoTConsensusConfig retryConfig = + IoTConsensusConfig.newBuilder() + .setReplication( + IoTConsensusConfig.Replication.newBuilder().setBasicRetryWaitTimeMs(10).build()) + .build(); + SyncStatus status = new SyncStatus(controller, retryConfig); + TLogEntry logEntry = new TLogEntry().setSearchIndex(1).setMemorySize(1); + Batch batch = new Batch(retryConfig); + batch.addTLogEntry(logEntry); + batch.buildIndex(); + + IoTConsensusMemoryManager memoryManager = IoTConsensusMemoryManager.getInstance(); + IMemoryBlock previousMemoryBlock = memoryManager.getMemoryBlock(); + CountDownLatch firstAllocationFailed = new CountDownLatch(1); + AtomicBoolean rejectAllocation = new AtomicBoolean(true); + IMemoryBlock memoryBlock = + new AtomicLongMemoryBlock("SyncStatusTest", null, batch.getMemorySize()) { + @Override + public boolean allocate(long sizeInByte) { + if (rejectAllocation.compareAndSet(true, false)) { + firstAllocationFailed.countDown(); + return false; + } + return super.allocate(sizeInByte); + } + }; + ExecutorService executor = Executors.newSingleThreadExecutor(); + memoryManager.setMemoryBlock(memoryBlock); + try { + Future<?> future = + executor.submit( + () -> { + status.addNextBatch(batch); + return null; + }); + + Assert.assertTrue(firstAllocationFailed.await(5, TimeUnit.SECONDS)); + future.get(5, TimeUnit.SECONDS); + + Assert.assertEquals(1, status.getPendingBatches().size()); + status.removeBatch(batch); + Assert.assertEquals(0, status.getPendingBatches().size()); + } finally { + executor.shutdownNow(); + executor.awaitTermination(5, TimeUnit.SECONDS); + status.free(); + memoryManager.setMemoryBlock(previousMemoryBlock); + } + } }
