This is an automated email from the ASF dual-hosted git repository.

Caideyipi pushed a commit to branch revert
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 988943b99c4a96f984ca18ee507772082e0b0e3a
Author: Caideyipi <[email protected]>
AuthorDate: Thu Aug 27 18:10:03 2026 +0800

    Revert "Fix IoTConsensus batch accumulation and config reload (#18501)"
    
    This reverts commit e591e66c0f522e5e8330874bcf5bbdfe63827388.
---
 .../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, 10 insertions(+), 288 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 477d8a5cb11..ae8cdd697c3 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 volatile IoTConsensusConfig config;
+  private IoTConsensusConfig config;
 
   /**
    * Optional callback invoked after a new local peer is created via {@link 
#createLocalPeer}. Used
@@ -539,9 +539,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 e074e7204ee..35070bb51b4 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
@@ -148,7 +148,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;
@@ -1354,7 +1354,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/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 39caeff33ed..6250b361e38 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,10 +412,6 @@ 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
@@ -446,7 +428,6 @@ public class LogDispatcher {
 
     public Batch getBatch() {
 
-      final IoTConsensusConfig currentConfig = config;
       long startIndex = syncStatus.getNextSendingIndex();
       long maxIndex;
       synchronized (impl.getIndexObject()) {
@@ -461,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();
@@ -475,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 3f4ea8ab528..00000000000
--- 
a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/logdispatcher/LogDispatcherTest.java
+++ /dev/null
@@ -1,224 +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 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