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 e591e66c0f5 Fix IoTConsensus batch accumulation and config reload
(#18501)
e591e66c0f5 is described below
commit e591e66c0f522e5e8330874bcf5bbdfe63827388
Author: Caideyipi <[email protected]>
AuthorDate: Mon Aug 24 11:06:39 2026 +0800
Fix IoTConsensus batch accumulation and config reload (#18501)
---
.../apache/iotdb/consensus/iot/IoTConsensus.java | 5 +-
.../consensus/iot/IoTConsensusServerImpl.java | 3 +-
.../logdispatcher/IoTConsensusMemoryManager.java | 6 +-
.../consensus/iot/logdispatcher/LogDispatcher.java | 31 ++-
.../consensus/iot/logdispatcher/SyncStatus.java | 7 +-
.../IoTConsensusMemoryManagerTest.java | 22 ++
.../iot/logdispatcher/LogDispatcherTest.java | 224 +++++++++++++++++++++
7 files changed, 288 insertions(+), 10 deletions(-)
diff --git
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensus.java
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensus.java
index 3cab3498b36..aa7ecbbf4ec 100644
---
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensus.java
+++
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensus.java
@@ -104,7 +104,7 @@ public class IoTConsensus implements IConsensus {
new ConcurrentHashMap<>();
private final IoTConsensusRPCService service;
private final RegisterManager registerManager = new RegisterManager();
- private IoTConsensusConfig config;
+ private volatile IoTConsensusConfig config;
/**
* Optional callback invoked after a new local peer is created via {@link
#createLocalPeer}. Used
@@ -537,6 +537,9 @@ public class IoTConsensus implements IConsensus {
public void reloadConsensusConfig(ConsensusConfig consensusConfig) {
config = consensusConfig.getIotConsensusConfig();
+ IoTConsensusMemoryManager.getInstance()
+
.updateMaxMemoryRatioForQueue(config.getReplication().getMaxMemoryRatioForQueue());
+
for (IoTConsensusServerImpl impl : stateMachineMap.values()) {
impl.reloadConsensusConfig(config);
}
diff --git
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
index 4cb109f0eb1..338bf30c1d9 100644
---
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
+++
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
@@ -147,7 +147,7 @@ public class IoTConsensusServerImpl {
private final Set<Peer> configuration = ConcurrentHashMap.newKeySet();
private final AtomicLong searchIndex;
private final LogDispatcher logDispatcher;
- private IoTConsensusConfig config;
+ private volatile IoTConsensusConfig config;
private final ConsensusReqReader consensusReqReader;
private volatile boolean active;
private String newSnapshotDirName;
@@ -1350,6 +1350,7 @@ public class IoTConsensusServerImpl {
/** This method is used for hot reload of IoTConsensusConfig. */
public void reloadConsensusConfig(IoTConsensusConfig config) {
this.config = config;
+ logDispatcher.reloadConfig(config);
}
/**
diff --git
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManager.java
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManager.java
index 161494a5fe8..1247a45129e 100644
---
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManager.java
+++
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManager.java
@@ -37,7 +37,7 @@ public class IoTConsensusMemoryManager {
private final AtomicLong syncMemorySizeInByte = new AtomicLong(0);
private IMemoryBlock memoryBlock =
new AtomicLongMemoryBlock("Consensus-Default", null,
Runtime.getRuntime().maxMemory() / 10);
- private Double maxMemoryRatioForQueue = 0.6;
+ private volatile double maxMemoryRatioForQueue = 0.6;
private IoTConsensusMemoryManager() {
MetricService.getInstance().addMetricSet(new
IoTConsensusMemoryManagerMetrics(this));
@@ -158,6 +158,10 @@ public class IoTConsensusMemoryManager {
this.maxMemoryRatioForQueue = maxMemoryRatioForQueue;
}
+ public void updateMaxMemoryRatioForQueue(double maxMemoryRatioForQueue) {
+ this.maxMemoryRatioForQueue = maxMemoryRatioForQueue;
+ }
+
@TestOnly
public void reset() {
this.memoryBlock.release(this.memoryBlock.getUsedMemoryInBytes());
diff --git
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
index 6250b361e38..39caeff33ed 100644
---
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
+++
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcher.java
@@ -183,6 +183,10 @@ public class LogDispatcher {
impl.checkAndUpdateSafeDeletedSearchIndex();
}
+ public synchronized void reloadConfig(IoTConsensusConfig config) {
+ threads.forEach(thread -> thread.reloadConfig(config));
+ }
+
public void offer(IndexedConsensusRequest request) {
offer(request, true);
}
@@ -230,7 +234,7 @@ public class LogDispatcher {
private static final long PENDING_REQUEST_TAKING_TIME_OUT_IN_MS = 10_000L;
private static final long START_INDEX = 1;
- private final IoTConsensusConfig config;
+ private volatile IoTConsensusConfig config;
private final Peer peer;
private final IndexController controller;
// A sliding window class that manages asynchronous pendingBatches
@@ -289,6 +293,11 @@ public class LogDispatcher {
return config;
}
+ private void reloadConfig(IoTConsensusConfig config) {
+ this.config = config;
+ syncStatus.reloadConfig(config);
+ }
+
public int getPendingEntriesSize() {
return pendingEntries.size();
}
@@ -375,11 +384,16 @@ public class LogDispatcher {
IndexedConsensusRequest request =
pendingEntries.poll(calculateIdlePollTimeoutInMs(),
TimeUnit.MILLISECONDS);
if (request != null) {
+ final IoTConsensusConfig currentConfig = config;
+ final boolean shouldWaitForBatchAccumulation =
+ pendingEntries.size()
+ <=
currentConfig.getReplication().getMaxLogEntriesNumPerBatch()
+ && bufferedEntries.isEmpty();
bufferedEntries.add(request);
// If write pressure is low, we simply sleep a little to reduce
the number of RPC
- if (pendingEntries.size() <=
config.getReplication().getMaxLogEntriesNumPerBatch()
- && bufferedEntries.isEmpty()) {
-
Thread.sleep(config.getReplication().getMaxWaitingTimeForAccumulatingBatchInMs());
+ if (shouldWaitForBatchAccumulation) {
+ waitForBatchAccumulation(
+
currentConfig.getReplication().getMaxWaitingTimeForAccumulatingBatchInMs());
}
} else {
maybeSendIdleWriterSafeTimeBarrier();
@@ -412,6 +426,10 @@ public class LogDispatcher {
logger.info(IoTConsensusMessages.DISPATCHER_EXITS, impl.getThisNode(),
peer);
}
+ void waitForBatchAccumulation(long waitingTimeInMs) throws
InterruptedException {
+ Thread.sleep(waitingTimeInMs);
+ }
+
public void updateSafelyDeletedSearchIndex() {
// update safely deleted search index to delete outdated info,
// indicating that insert nodes whose search index are before this value
can be deleted
@@ -428,6 +446,7 @@ public class LogDispatcher {
public Batch getBatch() {
+ final IoTConsensusConfig currentConfig = config;
long startIndex = syncStatus.getNextSendingIndex();
long maxIndex;
synchronized (impl.getIndexObject()) {
@@ -442,7 +461,7 @@ public class LogDispatcher {
// Use drainTo instead of poll to reduce lock overhead
pendingEntries.drainTo(
bufferedEntries,
- config.getReplication().getMaxLogEntriesNumPerBatch() -
bufferedEntries.size());
+ currentConfig.getReplication().getMaxLogEntriesNumPerBatch() -
bufferedEntries.size());
}
// remove all request that searchIndex < startIndex
Iterator<IndexedConsensusRequest> iterator = bufferedEntries.iterator();
@@ -456,7 +475,7 @@ public class LogDispatcher {
}
}
- Batch batches = new Batch(config);
+ Batch batches = new Batch(currentConfig);
// This condition will be executed in several scenarios:
// 1. restart
// 2. The getBatch() is invoked immediately at the moment the
PendingEntries are consumed
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 1749384f549..3df8a720614 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
@@ -32,7 +32,7 @@ import java.util.List;
public class SyncStatus {
private static final Logger LOGGER =
LoggerFactory.getLogger(SyncStatus.class);
- private final IoTConsensusConfig config;
+ private IoTConsensusConfig config;
private final IndexController controller;
private final LinkedList<Batch> pendingBatches = new LinkedList<>();
private final IoTConsensusMemoryManager iotConsensusMemoryManager =
@@ -43,6 +43,11 @@ public class SyncStatus {
this.config = config;
}
+ public synchronized void reloadConfig(IoTConsensusConfig config) {
+ this.config = config;
+ notifyAll();
+ }
+
/**
* we may block here if the synchronization pipeline is full.
*
diff --git
a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManagerTest.java
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManagerTest.java
index 9bddb298716..88ea61cb2af 100644
---
a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManagerTest.java
+++
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/IoTConsensusMemoryManagerTest.java
@@ -41,11 +41,14 @@ import static org.junit.Assert.assertTrue;
public class IoTConsensusMemoryManagerTest {
private IMemoryBlock previousMemoryBlock;
+ private double previousMaxMemoryRatioForQueue;
private long memoryBlockSize = 16 * 1024L;
@Before
public void setUp() throws Exception {
previousMemoryBlock =
IoTConsensusMemoryManager.getInstance().getMemoryBlock();
+ previousMaxMemoryRatioForQueue =
+ IoTConsensusMemoryManager.getInstance().getMaxMemoryRatioForQueue();
IoTConsensusMemoryManager.getInstance()
.setMemoryBlock(new AtomicLongMemoryBlock("Test", null,
memoryBlockSize));
IoTConsensusMemoryManager.getInstance().reset();
@@ -55,6 +58,8 @@ public class IoTConsensusMemoryManagerTest {
public void tearDown() throws Exception {
IoTConsensusMemoryManager.getInstance().reset();
IoTConsensusMemoryManager.getInstance().setMemoryBlock(previousMemoryBlock);
+ IoTConsensusMemoryManager.getInstance()
+ .updateMaxMemoryRatioForQueue(previousMaxMemoryRatioForQueue);
}
@Test
@@ -121,6 +126,23 @@ public class IoTConsensusMemoryManagerTest {
assertEquals(0L, request.getRetainedMemorySize());
}
+ @Test
+ public void testUpdateMaxMemoryRatioForQueue() {
+ final IndexedConsensusRequest request =
+ new IndexedConsensusRequest(
+ 1,
+ Collections.singletonList(
+ new ByteBufferConsensusRequest(ByteBuffer.allocate((int)
(memoryBlockSize / 3)))));
+ request.buildSerializedRequests();
+
+ IoTConsensusMemoryManager.getInstance().updateMaxMemoryRatioForQueue(0.25);
+ assertFalse(IoTConsensusMemoryManager.getInstance().reserve(request));
+
+ IoTConsensusMemoryManager.getInstance().updateMaxMemoryRatioForQueue(0.5);
+ assertTrue(IoTConsensusMemoryManager.getInstance().reserve(request));
+ IoTConsensusMemoryManager.getInstance().free(request);
+ }
+
private void testReserveAndRelease(int numReservation) {
int allocationSize = 1;
long allocatedSize = 0;
diff --git
a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcherTest.java
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcherTest.java
new file mode 100644
index 00000000000..3f4ea8ab528
--- /dev/null
+++
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcherTest.java
@@ -0,0 +1,224 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+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.disk.strategy.DirectoryStrategyType;
+import org.apache.iotdb.consensus.common.Peer;
+import org.apache.iotdb.consensus.common.request.IndexedConsensusRequest;
+import org.apache.iotdb.consensus.config.IoTConsensusConfig;
+import org.apache.iotdb.consensus.iot.IoTConsensusServerImpl;
+import org.apache.iotdb.consensus.iot.client.DispatchLogHandler;
+import org.apache.iotdb.consensus.iot.thrift.TLogEntry;
+import org.apache.iotdb.consensus.iot.util.TestEntry;
+import org.apache.iotdb.consensus.iot.util.TestStateMachine;
+
+import org.junit.Rule;
+import org.junit.Test;
+import org.junit.rules.TemporaryFolder;
+
+import java.lang.reflect.Field;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertSame;
+import static org.junit.Assert.assertTrue;
+
+public class LogDispatcherTest {
+
+ @Rule public final TemporaryFolder temporaryFolder = new TemporaryFolder();
+
+ @Test
+ public void testWaitForBatchAccumulationAfterFirstRequest() throws Exception
{
+ final Peer localPeer = createPeer(1, 6667);
+ final Peer remotePeer = createPeer(2, 6668);
+ final IoTConsensusConfig config = IoTConsensusConfig.newBuilder().build();
+ final ScheduledExecutorService backgroundTaskService =
+ Executors.newSingleThreadScheduledExecutor();
+ final ExecutorService executorService =
Executors.newSingleThreadExecutor();
+ LogDispatcher.LogDispatcherThread dispatcherThread = null;
+ Future<?> dispatcherFuture = null;
+ try {
+ final IoTConsensusServerImpl server =
+ createServer(
+ localPeer, Collections.singletonList(localPeer), config,
backgroundTaskService);
+ final Batch batch = createBatch(config, 1);
+ final CountDownLatch accumulationWaitInvoked = new CountDownLatch(1);
+ final AtomicInteger getBatchInvocations = new AtomicInteger();
+ dispatcherThread =
+ server.getLogDispatcher().new LogDispatcherThread(remotePeer,
config, 0) {
+ @Override
+ public Batch getBatch() {
+ return getBatchInvocations.getAndIncrement() == 0 ? new
Batch(config) : batch;
+ }
+
+ @Override
+ void waitForBatchAccumulation(long waitingTimeInMs) {
+ accumulationWaitInvoked.countDown();
+ }
+
+ @Override
+ public void sendBatchAsync(Batch sentBatch, DispatchLogHandler
handler) {
+ getSyncStatus().removeBatch(sentBatch);
+ Thread.currentThread().interrupt();
+ }
+ };
+ assertTrue(
+ dispatcherThread.offer(
+ new IndexedConsensusRequest(
+ 1, Collections.singletonList(new TestEntry(1, localPeer)))));
+
+ dispatcherFuture = executorService.submit(dispatcherThread);
+
+ assertTrue(accumulationWaitInvoked.await(5, TimeUnit.SECONDS));
+ dispatcherFuture.get(5, TimeUnit.SECONDS);
+ } finally {
+ if (dispatcherFuture != null) {
+ dispatcherFuture.cancel(true);
+ }
+ executorService.shutdownNow();
+ executorService.awaitTermination(5, TimeUnit.SECONDS);
+ if (dispatcherThread != null) {
+ dispatcherThread.stop();
+ }
+ backgroundTaskService.shutdownNow();
+ }
+ }
+
+ @Test
+ public void testReloadConfigUpdatesExistingDispatcherPipeline() throws
Exception {
+ final Peer localPeer = createPeer(1, 6677);
+ final Peer remotePeer = createPeer(2, 6678);
+ final IoTConsensusConfig initialConfig =
+ IoTConsensusConfig.newBuilder()
+ .setReplication(
+ IoTConsensusConfig.Replication.newBuilder()
+ .setMaxLogEntriesNumPerBatch(1)
+ .setMaxPendingBatchesNum(1)
+ .build())
+ .build();
+ final ScheduledExecutorService backgroundTaskService =
+ Executors.newSingleThreadScheduledExecutor();
+ final ExecutorService executorService =
Executors.newSingleThreadExecutor();
+ LogDispatcher dispatcher = null;
+ Future<?> secondBatchFuture = null;
+ try {
+ final IoTConsensusServerImpl server =
+ createServer(
+ localPeer,
+ Arrays.asList(localPeer, remotePeer),
+ initialConfig,
+ backgroundTaskService);
+ dispatcher = server.getLogDispatcher();
+ final LogDispatcher.LogDispatcherThread dispatcherThread =
getOnlyThread(dispatcher);
+ dispatcher.start();
+
+ final SyncStatus syncStatus = dispatcherThread.getSyncStatus();
+ syncStatus.addNextBatch(createBatch(initialConfig, 1));
+ final CountDownLatch secondBatchAttempted = new CountDownLatch(1);
+ secondBatchFuture =
+ executorService.submit(
+ () -> {
+ secondBatchAttempted.countDown();
+ syncStatus.addNextBatch(createBatch(initialConfig, 2));
+ return null;
+ });
+ assertTrue(secondBatchAttempted.await(5, TimeUnit.SECONDS));
+ Thread.sleep(100);
+ assertFalse(secondBatchFuture.isDone());
+
+ final IoTConsensusConfig reloadedConfig =
+ IoTConsensusConfig.newBuilder()
+ .setReplication(
+ IoTConsensusConfig.Replication.newBuilder()
+ .setMaxLogEntriesNumPerBatch(2)
+ .setMaxPendingBatchesNum(2)
+ .build())
+ .build();
+ server.reloadConsensusConfig(reloadedConfig);
+
+ secondBatchFuture.get(5, TimeUnit.SECONDS);
+ assertSame(reloadedConfig, dispatcherThread.getConfig());
+ assertEquals(2, syncStatus.getPendingBatches().size());
+ } finally {
+ if (secondBatchFuture != null) {
+ secondBatchFuture.cancel(true);
+ }
+ executorService.shutdownNow();
+ executorService.awaitTermination(5, TimeUnit.SECONDS);
+ if (dispatcher != null) {
+ dispatcher.stop();
+ }
+ backgroundTaskService.shutdownNow();
+ }
+ }
+
+ private IoTConsensusServerImpl createServer(
+ Peer localPeer,
+ List<Peer> configuration,
+ IoTConsensusConfig config,
+ ScheduledExecutorService backgroundTaskService)
+ throws Exception {
+ return new IoTConsensusServerImpl(
+ temporaryFolder.newFolder().getAbsolutePath(),
+ null,
+ DirectoryStrategyType.SEQUENCE_STRATEGY,
+ localPeer,
+ configuration,
+ new TestStateMachine(),
+ backgroundTaskService,
+ null,
+ null,
+ config);
+ }
+
+ private static Peer createPeer(int nodeId, int port) {
+ return new Peer(new DataRegionId(1), nodeId, new TEndPoint("127.0.0.1",
port));
+ }
+
+ private static Batch createBatch(IoTConsensusConfig config, long
searchIndex) {
+ final Batch batch = new Batch(config);
+ batch.addTLogEntry(new
TLogEntry().setSearchIndex(searchIndex).setMemorySize(1));
+ batch.buildIndex();
+ return batch;
+ }
+
+ @SuppressWarnings("unchecked")
+ private static LogDispatcher.LogDispatcherThread getOnlyThread(LogDispatcher
dispatcher)
+ throws Exception {
+ final Field threadsField = LogDispatcher.class.getDeclaredField("threads");
+ threadsField.setAccessible(true);
+ final List<LogDispatcher.LogDispatcherThread> threads =
+ (List<LogDispatcher.LogDispatcherThread>) threadsField.get(dispatcher);
+ assertEquals(1, threads.size());
+ return threads.get(0);
+ }
+}