From a900d5062a5d390900f01213f772c084d8afdb5d Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Fri, 21 Aug 2026 12:26:05 +0800 Subject: [PATCH] Fix IoTConsensus batch accumulation and config reload --- .../iotdb/consensus/iot/IoTConsensus.java | 5 +- .../consensus/iot/IoTConsensusServerImpl.java | 3 +- .../IoTConsensusMemoryManager.java | 6 +- .../iot/logdispatcher/LogDispatcher.java | 31 ++- .../iot/logdispatcher/SyncStatus.java | 7 +- .../IoTConsensusMemoryManagerTest.java | 22 ++ .../iot/logdispatcher/LogDispatcherTest.java | 224 ++++++++++++++++++ 7 files changed, 288 insertions(+), 10 deletions(-) create mode 100644 iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcherTest.java 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 3cab3498b36a0..aa7ecbbf4ec94 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 String getRegionDirFromConsensusGroupId(ConsensusGroupId groupId) { 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 4cb109f0eb137..338bf30c1d961 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 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 String getConsensusGroupId() { /** 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 161494a5fe8e4..1247a45129e1a 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 void init(IMemoryBlock memoryBlock, double maxMemoryRatioForQueue) { 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 6250b361e389d..39caeff33ed55 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 synchronized void checkAndFlushIndex() { 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 LogDispatcherThread implements Runnable { 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 IoTConsensusConfig getConfig() { return config; } + private void reloadConfig(IoTConsensusConfig config) { + this.config = config; + syncStatus.reloadConfig(config); + } + public int getPendingEntriesSize() { return pendingEntries.size(); } @@ -375,11 +384,16 @@ public void run() { 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 void run() { 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 void updateSafelyDeletedSearchIndex() { public Batch getBatch() { + final IoTConsensusConfig currentConfig = config; long startIndex = syncStatus.getNextSendingIndex(); long maxIndex; synchronized (impl.getIndexObject()) { @@ -442,7 +461,7 @@ public Batch getBatch() { // 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 iterator = bufferedEntries.iterator(); @@ -456,7 +475,7 @@ public Batch getBatch() { } } - 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 1749384f54911..3df8a720614ac 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 @@ 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 pendingBatches = new LinkedList<>(); private final IoTConsensusMemoryManager iotConsensusMemoryManager = @@ -43,6 +43,11 @@ public SyncStatus(IndexController controller, IoTConsensusConfig config) { 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 9bddb29871617..88ea61cb2af61 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 @@ 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 void setUp() throws Exception { public void tearDown() throws Exception { IoTConsensusMemoryManager.getInstance().reset(); IoTConsensusMemoryManager.getInstance().setMemoryBlock(previousMemoryBlock); + IoTConsensusMemoryManager.getInstance() + .updateMaxMemoryRatioForQueue(previousMaxMemoryRatioForQueue); } @Test @@ -121,6 +126,23 @@ public void testClearUnserializedRequest() { 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 0000000000000..3f4ea8ab528a5 --- /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 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 threads = + (List) threadsField.get(dispatcher); + assertEquals(1, threads.size()); + return threads.get(0); + } +}