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

Reply via email to