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 448a805c7b6 Revert IoTConsensus batch accumulation changes (#18598)
448a805c7b6 is described below

commit 448a805c7b61433addc44deefa0861466e20ebd4
Author: Caideyipi <[email protected]>
AuthorDate: Tue Sep 8 16:49:13 2026 +0800

    Revert IoTConsensus batch accumulation changes (#18598)
    
    * Revert "Fix IoTConsensus batch accumulation by byte size (#18560)"
    
    This reverts commit c41838e1386706003125f333ddadb25f6132433d.
    
    * Revert "Fix IoTConsensus batch accumulation latency (#18522)"
    
    This reverts commit ebd9c2072e3195a510d784e7b0fb4422fa419935.
---
 .../apache/iotdb/consensus/iot/IoTConsensus.java   |   5 +-
 .../consensus/iot/IoTConsensusServerImpl.java      |   3 +-
 .../iotdb/consensus/iot/logdispatcher/Batch.java   |   6 +-
 .../logdispatcher/IoTConsensusMemoryManager.java   |   6 +-
 .../consensus/iot/logdispatcher/LogDispatcher.java |  59 +---
 .../consensus/iot/logdispatcher/SyncStatus.java    |   7 +-
 .../IoTConsensusMemoryManagerTest.java             |  22 --
 .../iot/logdispatcher/LogDispatcherTest.java       | 354 ---------------------
 8 files changed, 11 insertions(+), 451 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 ede52c92ed3..9c0692132f0 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
@@ -108,7 +108,7 @@ public class IoTConsensus implements IConsensus {
   private final RegisterManager registerManager = new RegisterManager();
   private final UserDataTransferAuditHandler userDataTransferAuditHandler;
   private final UserDataTransferAuditClassifier 
userDataTransferAuditClassifier;
-  private volatile IoTConsensusConfig config;
+  private IoTConsensusConfig config;
 
   /**
    * Optional callback invoked after a new local peer is created via {@link 
#createLocalPeer}. Used
@@ -549,9 +549,6 @@ 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 6e86ec72a3e..f4448e81f34 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
@@ -152,7 +152,7 @@ public class IoTConsensusServerImpl {
   private final Set<Peer> configuration = ConcurrentHashMap.newKeySet();
   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;
@@ -1474,7 +1474,6 @@ 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/Batch.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/logdispatcher/Batch.java
index 26fcddb5d2a..55569b8a34f 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
@@ -69,10 +69,6 @@ public class Batch {
   }
 
   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.
@@ -81,7 +77,7 @@ public class Batch {
     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 1247a45129e..161494a5fe8 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 volatile double maxMemoryRatioForQueue = 0.6;
+  private Double maxMemoryRatioForQueue = 0.6;
 
   private IoTConsensusMemoryManager() {
     MetricService.getInstance().addMetricSet(new 
IoTConsensusMemoryManagerMetrics(this));
@@ -158,10 +158,6 @@ 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 d9b34b9f79d..304badbd8ae 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,10 +183,6 @@ 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);
   }
@@ -234,7 +230,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 volatile IoTConsensusConfig config;
+    private final IoTConsensusConfig config;
     private final Peer peer;
     private final IndexController controller;
     // A sliding window class that manages asynchronous pendingBatches
@@ -293,11 +289,6 @@ public class LogDispatcher {
       return config;
     }
 
-    private void reloadConfig(IoTConsensusConfig config) {
-      this.config = config;
-      syncStatus.reloadConfig(config);
-    }
-
     public int getPendingEntriesSize() {
       return pendingEntries.size();
     }
@@ -384,16 +375,11 @@ 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 (shouldWaitForBatchAccumulation) {
-                waitForBatchAccumulation(
-                    
currentConfig.getReplication().getMaxWaitingTimeForAccumulatingBatchInMs());
+              if (pendingEntries.size() <= 
config.getReplication().getMaxLogEntriesNumPerBatch()
+                  && bufferedEntries.isEmpty()) {
+                
Thread.sleep(config.getReplication().getMaxWaitingTimeForAccumulatingBatchInMs());
               }
             } else {
               maybeSendIdleWriterSafeTimeBarrier();
@@ -426,38 +412,6 @@ public class LogDispatcher {
       logger.info(IoTConsensusMessages.DISPATCHER_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
@@ -474,7 +428,6 @@ public class LogDispatcher {
 
     public Batch getBatch() {
 
-      final IoTConsensusConfig currentConfig = config;
       long startIndex = syncStatus.getNextSendingIndex();
       long maxIndex;
       synchronized (impl.getIndexObject()) {
@@ -489,7 +442,7 @@ public class LogDispatcher {
         // 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<IndexedConsensusRequest> iterator = bufferedEntries.iterator();
@@ -503,7 +456,7 @@ public class LogDispatcher {
         }
       }
 
-      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 3df8a720614..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
@@ -32,7 +32,7 @@ import java.util.List;
 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<Batch> pendingBatches = new LinkedList<>();
   private final IoTConsensusMemoryManager iotConsensusMemoryManager =
@@ -43,11 +43,6 @@ 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 88ea61cb2af..9bddb298716 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,14 +41,11 @@ 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();
@@ -58,8 +55,6 @@ public class IoTConsensusMemoryManagerTest {
   public void tearDown() throws Exception {
     IoTConsensusMemoryManager.getInstance().reset();
     
IoTConsensusMemoryManager.getInstance().setMemoryBlock(previousMemoryBlock);
-    IoTConsensusMemoryManager.getInstance()
-        .updateMaxMemoryRatioForQueue(previousMaxMemoryRatioForQueue);
   }
 
   @Test
@@ -126,23 +121,6 @@ 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
deleted file mode 100644
index a6e299748d2..00000000000
--- 
a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcherTest.java
+++ /dev/null
@@ -1,354 +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.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 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<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);
-  }
-}

Reply via email to