This is an automated email from the ASF dual-hosted git repository.
Caideyipi 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 a10bc6ba7e7 Revert "Fix IoTConsensus batch accumulation and config
reload (#18501)" (#18535)
a10bc6ba7e7 is described below
commit a10bc6ba7e744eb3981506613b2c439a0abacba3
Author: Caideyipi <[email protected]>
AuthorDate: Thu Aug 27 18:11:41 2026 +0800
Revert "Fix IoTConsensus batch accumulation and config reload (#18501)"
(#18535)
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);
- }
-}