From fdf65096c2e6d9564b940098e1fd8108d570c2e6 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Wed, 9 Sep 2026 10:42:59 +0800 Subject: [PATCH] Revert IoTConsensus batch accumulation changes (#18598) --- .../iotdb/consensus/iot/IoTConsensus.java | 7 +- .../consensus/iot/IoTConsensusServerImpl.java | 3 +- .../consensus/iot/logdispatcher/Batch.java | 6 +- .../IoTConsensusMemoryManager.java | 4 +- .../iot/logdispatcher/LogDispatcher.java | 59 +-- .../iot/logdispatcher/SyncStatus.java | 7 +- .../IoTConsensusMemoryManagerTest.java | 14 - .../iot/logdispatcher/LogDispatcherTest.java | 352 ------------------ 8 files changed, 12 insertions(+), 440 deletions(-) delete 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 11707fe5634db..d15d6e365a77f 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 @@ -97,7 +97,7 @@ public class IoTConsensus implements IConsensus { new ConcurrentHashMap<>(); private final IoTConsensusRPCService service; private final RegisterManager registerManager = new RegisterManager(); - private volatile IoTConsensusConfig config; + private IoTConsensusConfig config; private final IClientManager clientManager; private final IClientManager syncClientManager; private final ScheduledExecutorService backgroundTaskService; @@ -472,11 +472,6 @@ public String getRegionDirFromConsensusGroupId(ConsensusGroupId groupId) { public void reloadConsensusConfig(ConsensusConfig consensusConfig) { config = consensusConfig.getIotConsensusConfig(); - IoTConsensusMemoryManager.getInstance() - .init( - config.getReplication().getAllocateMemoryForConsensus(), - config.getReplication().getAllocateMemoryForQueue()); - 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 7033ddf36d960..3002b018e3ea8 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 @@ -116,7 +116,7 @@ public class IoTConsensusServerImpl { private final TreeSet configuration; private final AtomicLong searchIndex; private final LogDispatcher logDispatcher; - private volatile IoTConsensusConfig config; + private IoTConsensusConfig config; private final ConsensusReqReader consensusReqReader; private volatile boolean active; private String newSnapshotDirName; @@ -911,7 +911,6 @@ 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/Batch.java b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/Batch.java index 8d28743f7ba23..72b68ab96ac7e 100644 --- a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/Batch.java +++ b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/Batch.java @@ -63,10 +63,6 @@ public void addTLogEntry(TLogEntry entry) { } public boolean canAccumulate() { - return canAccumulate(config, logEntries.size(), memorySize); - } - - static boolean canAccumulate(IoTConsensusConfig config, int logEntriesSize, long memorySize) { // When reading entries from the WAL, the memory size is calculated based on the serialized // size, which can be significantly smaller than the actual size. // Thus, we add a multiplier to sender's memory size to estimate the receiver's memory cost. @@ -75,7 +71,7 @@ static boolean canAccumulate(IoTConsensusConfig config, int logEntriesSize, long long senderMemSize = LogDispatcher.getSenderMemSizeSum().get(); double multiplier = senderMemSize > 0 ? (double) receiverMemSize / senderMemSize : 1.0; multiplier = Math.max(multiplier, 1.0); - return logEntriesSize < config.getReplication().getMaxLogEntriesNumPerBatch() + return logEntries.size() < config.getReplication().getMaxLogEntriesNumPerBatch() && ((long) (memorySize * multiplier)) < config.getReplication().getMaxSizePerBatch(); } 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 d8adec09a7b06..22e5484f5a203 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 @@ -33,8 +33,8 @@ public class IoTConsensusMemoryManager { private final AtomicLong memorySizeInByte = new AtomicLong(0); private final AtomicLong queueMemorySizeInByte = new AtomicLong(0); private final AtomicLong syncMemorySizeInByte = new AtomicLong(0); - private volatile long maxMemorySizeInByte = Runtime.getRuntime().maxMemory() / 10; - private volatile long maxMemorySizeForQueueInByte = Runtime.getRuntime().maxMemory() / 100 * 6; + private Long maxMemorySizeInByte = Runtime.getRuntime().maxMemory() / 10; + private Long maxMemorySizeForQueueInByte = Runtime.getRuntime().maxMemory() / 100 * 6; private IoTConsensusMemoryManager() { MetricService.getInstance().addMetricSet(new IoTConsensusMemoryManagerMetrics(this)); 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 04d8371977294..374691bf38bf1 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 @@ -179,10 +179,6 @@ public void checkAndFlushIndex() { } } - public synchronized void reloadConfig(IoTConsensusConfig config) { - threads.forEach(thread -> thread.reloadConfig(config)); - } - public void offer(IndexedConsensusRequest request) { // we don't need to serialize and offer request when replicaNum is 1. if (!threads.isEmpty()) { @@ -219,7 +215,7 @@ public class LogDispatcherThread implements Runnable { private static final long PENDING_REQUEST_TAKING_TIME_OUT_IN_SEC = 10; private static final long START_INDEX = 1; - private volatile IoTConsensusConfig config; + private final IoTConsensusConfig config; private final Peer peer; private final IndexController controller; // A sliding window class that manages asynchronous pendingBatches @@ -277,11 +273,6 @@ public IoTConsensusConfig getConfig() { return config; } - private void reloadConfig(IoTConsensusConfig config) { - this.config = config; - syncStatus.reloadConfig(config); - } - public int getPendingEntriesSize() { return pendingEntries.size(); } @@ -367,16 +358,11 @@ public void run() { IndexedConsensusRequest request = pendingEntries.poll(PENDING_REQUEST_TAKING_TIME_OUT_IN_SEC, TimeUnit.SECONDS); 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 (shouldWaitForBatchAccumulation) { - waitForBatchAccumulation( - currentConfig.getReplication().getMaxWaitingTimeForAccumulatingBatchInMs()); + if (pendingEntries.size() <= config.getReplication().getMaxLogEntriesNumPerBatch() + && bufferedEntries.isEmpty()) { + Thread.sleep(config.getReplication().getMaxWaitingTimeForAccumulatingBatchInMs()); } } // Immediately check for interrupts after poll and sleep @@ -406,38 +392,6 @@ public void run() { logger.info("{}: Dispatcher for {} exits", impl.getThisNode(), peer); } - void waitForBatchAccumulation(long waitingTimeInMs) throws InterruptedException { - if (waitingTimeInMs <= 0) { - return; - } - - final long deadlineNanos = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(waitingTimeInMs); - final IoTConsensusConfig currentConfig = config; - int accumulatedEntries = bufferedEntries.size(); - long accumulatedMemorySize = - bufferedEntries.stream().mapToLong(IndexedConsensusRequest::getMemorySize).sum(); - - // Keep collecting while the batch is below both its entry and memory limits. A plain sleep, - // or checking only the entry limit, makes the dispatcher wait for the full accumulation - // interval after a batch has already reached its memory limit. This unnecessarily throttles - // IoTConsensus when each request contains a large tablet. - while (Batch.canAccumulate(currentConfig, accumulatedEntries, accumulatedMemorySize)) { - final long remainingNanos = deadlineNanos - System.nanoTime(); - if (remainingNanos <= 0) { - return; - } - - final IndexedConsensusRequest request = - pendingEntries.poll(remainingNanos, TimeUnit.NANOSECONDS); - if (request == null) { - return; - } - bufferedEntries.add(request); - accumulatedEntries++; - accumulatedMemorySize += request.getMemorySize(); - } - } - 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 @@ -452,7 +406,6 @@ public void updateSafelyDeletedSearchIndex() { } public Batch getBatch() { - final IoTConsensusConfig currentConfig = config; long startIndex = syncStatus.getNextSendingIndex(); long maxIndex; synchronized (impl.getIndexObject()) { @@ -467,7 +420,7 @@ public Batch getBatch() { // Use drainTo instead of poll to reduce lock overhead pendingEntries.drainTo( bufferedEntries, - currentConfig.getReplication().getMaxLogEntriesNumPerBatch() - bufferedEntries.size()); + config.getReplication().getMaxLogEntriesNumPerBatch() - bufferedEntries.size()); } // remove all request that searchIndex < startIndex Iterator iterator = bufferedEntries.iterator(); @@ -481,7 +434,7 @@ public Batch getBatch() { } } - Batch batches = new Batch(currentConfig); + Batch batches = new Batch(config); // 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 9e6af37582845..506957c2b6864 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 @@ -31,7 +31,7 @@ public class SyncStatus { private static final Logger LOGGER = LoggerFactory.getLogger(SyncStatus.class); - private IoTConsensusConfig config; + private final IoTConsensusConfig config; private final IndexController controller; private final LinkedList pendingBatches = new LinkedList<>(); private final IoTConsensusMemoryManager iotConsensusMemoryManager = @@ -42,11 +42,6 @@ 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 6d6bae6165e1d..f87d8cd7f9887 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 @@ -37,20 +37,6 @@ public class IoTConsensusMemoryManagerTest { - @Test - public void testInitUpdatesMemoryLimits() { - IoTConsensusMemoryManager memoryManager = IoTConsensusMemoryManager.getInstance(); - long previousMaxMemory = memoryManager.getMaxMemorySizeInByte(); - long previousMaxQueueMemory = memoryManager.getMaxMemorySizeForQueueInByte(); - try { - memoryManager.init(1024, 512); - assertEquals(1024L, memoryManager.getMaxMemorySizeInByte().longValue()); - assertEquals(512L, memoryManager.getMaxMemorySizeForQueueInByte().longValue()); - } finally { - memoryManager.init(previousMaxMemory, previousMaxQueueMemory); - } - } - @Test public void testAllocateQueue() { IoTConsensusMemoryManager memoryManager = IoTConsensusMemoryManager.getInstance(); 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 deleted file mode 100644 index 7850785b77480..0000000000000 --- a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcherTest.java +++ /dev/null @@ -1,352 +0,0 @@ -/* - * 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.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.TreeSet; -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 testBatchAccumulationStopsWhenBatchIsFull() throws Exception { - final Peer localPeer = createPeer(1, 6687); - final Peer remotePeer = createPeer(2, 6688); - final IoTConsensusConfig config = - IoTConsensusConfig.newBuilder() - .setReplication( - IoTConsensusConfig.Replication.newBuilder() - .setMaxLogEntriesNumPerBatch(2) - .setMaxWaitingTimeForAccumulatingBatchInMs(10_000) - .build()) - .build(); - final ScheduledExecutorService backgroundTaskService = - Executors.newSingleThreadScheduledExecutor(); - final ExecutorService executorService = Executors.newSingleThreadExecutor(); - LogDispatcher.LogDispatcherThread dispatcherThread = null; - Future dispatcherFuture = null; - try { - final IoTConsensusServerImpl server = - createServer( - localPeer, Arrays.asList(localPeer, remotePeer), config, backgroundTaskService); - final CountDownLatch batchSent = 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) - : createBatch(config, 1); - } - - @Override - public void sendBatchAsync(Batch sentBatch, DispatchLogHandler handler) { - assertEquals(0, getPendingEntriesSize()); - batchSent.countDown(); - Thread.currentThread().interrupt(); - } - }; - assertTrue( - dispatcherThread.offer( - new IndexedConsensusRequest( - 1, Collections.singletonList(new TestEntry(1, localPeer))))); - assertTrue( - dispatcherThread.offer( - new IndexedConsensusRequest( - 2, Collections.singletonList(new TestEntry(2, localPeer))))); - - dispatcherFuture = executorService.submit(dispatcherThread); - assertTrue(batchSent.await(2, TimeUnit.SECONDS)); - dispatcherFuture.get(2, 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 testBatchAccumulationStopsWhenMemoryLimitIsReached() throws Exception { - final Peer localPeer = createPeer(1, 6697); - final Peer remotePeer = createPeer(2, 6698); - final IoTConsensusConfig config = - IoTConsensusConfig.newBuilder() - .setReplication( - IoTConsensusConfig.Replication.newBuilder() - .setMaxLogEntriesNumPerBatch(1024) - .setMaxSizePerBatch(1) - .setMaxWaitingTimeForAccumulatingBatchInMs(10_000) - .build()) - .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 CountDownLatch batchSent = 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) - : createBatch(config, 1); - } - - @Override - public void sendBatchAsync(Batch sentBatch, DispatchLogHandler handler) { - assertEquals(1, getPendingEntriesSize()); - assertEquals(1, getBufferRequestSize()); - batchSent.countDown(); - Thread.currentThread().interrupt(); - } - }; - final IndexedConsensusRequest firstRequest = - new IndexedConsensusRequest(1, Collections.singletonList(new TestEntry(1, localPeer))); - firstRequest.buildSerializedRequests(); - final IndexedConsensusRequest secondRequest = - new IndexedConsensusRequest(2, Collections.singletonList(new TestEntry(2, localPeer))); - secondRequest.buildSerializedRequests(); - assertTrue(dispatcherThread.offer(firstRequest)); - assertTrue(dispatcherThread.offer(secondRequest)); - - dispatcherFuture = executorService.submit(dispatcherThread); - assertTrue(batchSent.await(2, TimeUnit.SECONDS)); - dispatcherFuture.get(2, 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(), - localPeer, - new TreeSet<>(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); - } -}