This is an automated email from the ASF dual-hosted git repository.
jt2594838 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 e509114c515 fix: retry SyncStatus batch memory reservation (#18453)
e509114c515 is described below
commit e509114c515602585da7cbcffb8daf93a297c90f
Author: Jiang Tian <[email protected]>
AuthorDate: Wed Aug 12 18:27:58 2026 +0800
fix: retry SyncStatus batch memory reservation (#18453)
---
.../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);
+ }
+ }
}