Skip to content

Commit c33ef8f

Browse files
jt2594838JackieTien97
authored andcommitted
fix: retry SyncStatus batch memory reservation (#18453)
1 parent b8265e3 commit c33ef8f

2 files changed

Lines changed: 73 additions & 4 deletions

File tree

‎iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/SyncStatus.java‎

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -49,10 +49,15 @@ public SyncStatus(IndexController controller, IoTConsensusConfig config) {
4949
* @throws InterruptedException
5050
*/
5151
public synchronized void addNextBatch(Batch batch) throws InterruptedException {
52-
while ((pendingBatches.size() >= config.getReplication().getMaxPendingBatchesNum()
53-
|| !iotConsensusMemoryManager.reserve(batch))
54-
&& !Thread.interrupted()) {
55-
wait();
52+
while (true) {
53+
while (pendingBatches.size() >= config.getReplication().getMaxPendingBatchesNum()) {
54+
wait();
55+
}
56+
if (iotConsensusMemoryManager.reserve(batch)) {
57+
break;
58+
}
59+
// Memory may be freed by another SyncStatus, which cannot notify this monitor.
60+
wait(Math.max(1, config.getReplication().getBasicRetryWaitTimeMs()));
5661
}
5762
if (LOGGER.isDebugEnabled()) {
5863
LOGGER.debug(

‎iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/SyncStatusTest.java‎

Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,8 @@
2121

2222
import org.apache.iotdb.common.rpc.thrift.TEndPoint;
2323
import org.apache.iotdb.commons.consensus.DataRegionId;
24+
import org.apache.iotdb.commons.memory.AtomicLongMemoryBlock;
25+
import org.apache.iotdb.commons.memory.IMemoryBlock;
2426
import org.apache.iotdb.consensus.common.Peer;
2527
import org.apache.iotdb.consensus.config.IoTConsensusConfig;
2628
import org.apache.iotdb.consensus.iot.thrift.TLogEntry;
@@ -36,7 +38,14 @@
3638
import java.util.ArrayList;
3739
import java.util.List;
3840
import java.util.concurrent.CompletableFuture;
41+
import java.util.concurrent.CountDownLatch;
3942
import java.util.concurrent.ExecutionException;
43+
import java.util.concurrent.ExecutorService;
44+
import java.util.concurrent.Executors;
45+
import java.util.concurrent.Future;
46+
import java.util.concurrent.TimeUnit;
47+
import java.util.concurrent.TimeoutException;
48+
import java.util.concurrent.atomic.AtomicBoolean;
4049

4150
public class SyncStatusTest {
4251

@@ -242,4 +251,59 @@ public void waitTest() throws InterruptedException, ExecutionException {
242251
Assert.assertEquals(
243252
config.getReplication().getMaxPendingBatchesNum() + 1, status.getNextSendingIndex());
244253
}
254+
255+
@Test
256+
public void testFirstBatchRetriesMemoryReservation()
257+
throws InterruptedException, ExecutionException, TimeoutException {
258+
IndexController controller =
259+
new IndexController(storageDir.getAbsolutePath(), peer, 0, CHECK_POINT_GAP);
260+
IoTConsensusConfig retryConfig =
261+
IoTConsensusConfig.newBuilder()
262+
.setReplication(
263+
IoTConsensusConfig.Replication.newBuilder().setBasicRetryWaitTimeMs(10).build())
264+
.build();
265+
SyncStatus status = new SyncStatus(controller, retryConfig);
266+
TLogEntry logEntry = new TLogEntry().setSearchIndex(1).setMemorySize(1);
267+
Batch batch = new Batch(retryConfig);
268+
batch.addTLogEntry(logEntry);
269+
batch.buildIndex();
270+
271+
IoTConsensusMemoryManager memoryManager = IoTConsensusMemoryManager.getInstance();
272+
IMemoryBlock previousMemoryBlock = memoryManager.getMemoryBlock();
273+
CountDownLatch firstAllocationFailed = new CountDownLatch(1);
274+
AtomicBoolean rejectAllocation = new AtomicBoolean(true);
275+
IMemoryBlock memoryBlock =
276+
new AtomicLongMemoryBlock("SyncStatusTest", null, batch.getMemorySize()) {
277+
@Override
278+
public boolean allocate(long sizeInByte) {
279+
if (rejectAllocation.compareAndSet(true, false)) {
280+
firstAllocationFailed.countDown();
281+
return false;
282+
}
283+
return super.allocate(sizeInByte);
284+
}
285+
};
286+
ExecutorService executor = Executors.newSingleThreadExecutor();
287+
memoryManager.setMemoryBlock(memoryBlock);
288+
try {
289+
Future<?> future =
290+
executor.submit(
291+
() -> {
292+
status.addNextBatch(batch);
293+
return null;
294+
});
295+
296+
Assert.assertTrue(firstAllocationFailed.await(5, TimeUnit.SECONDS));
297+
future.get(5, TimeUnit.SECONDS);
298+
299+
Assert.assertEquals(1, status.getPendingBatches().size());
300+
status.removeBatch(batch);
301+
Assert.assertEquals(0, status.getPendingBatches().size());
302+
} finally {
303+
executor.shutdownNow();
304+
executor.awaitTermination(5, TimeUnit.SECONDS);
305+
status.free();
306+
memoryManager.setMemoryBlock(previousMemoryBlock);
307+
}
308+
}
245309
}

0 commit comments

Comments
 (0)