diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/ScmConfigKeys.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/ScmConfigKeys.java
index 598d79a011c4..365a0ab16d37 100644
--- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/ScmConfigKeys.java
+++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/ScmConfigKeys.java
@@ -630,6 +630,10 @@ public final class ScmConfigKeys {
public static final long
OZONE_SCM_HA_DBTRANSACTIONBUFFER_FLUSH_INTERVAL_DEFAULT = 60 * 1000L;
+ public static final String OZONE_SCM_HA_DBTRANSACTIONBUFFER_FLUSH_PENDING_LIMIT =
+ "ozone.scm.ha.dbtransactionbuffer.flush.pending.limit";
+ public static final long OZONE_SCM_HA_DBTRANSACTIONBUFFER_FLUSH_PENDING_LIMIT_DEFAULT = 10000L;
+
public static final String NET_TOPOLOGY_NODE_SWITCH_MAPPING_IMPL_KEY =
"net.topology.node.switch.mapping.impl";
public static final String HDDS_CONTAINER_RATIS_STATEMACHINE_WRITE_WAIT_INTERVAL
diff --git a/hadoop-hdds/common/src/main/resources/ozone-default.xml b/hadoop-hdds/common/src/main/resources/ozone-default.xml
index 51e66a3618c0..c32d06434b03 100644
--- a/hadoop-hdds/common/src/main/resources/ozone-default.xml
+++ b/hadoop-hdds/common/src/main/resources/ozone-default.xml
@@ -4389,7 +4389,16 @@
SCM, OZONE
Wait duration for flush of buffered transaction.
-
+
+ ozone.scm.ha.dbtransactionbuffer.flush.pending.limit
+ 10000
+ SCM, OZONE
+
+ Maximum number of buffered DB operations the SCM HA transaction buffer holds
+ before applyTransaction flushes them to RocksDB. This bounds memory use while
+ a restarted follower applies a large backlog of committed Ratis log entries.
+
+
ozone.s3g.kerberos.keytab.file
diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMHADBTransactionBuffer.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMHADBTransactionBuffer.java
index 39f8c5ecda4f..097eb1d0a93c 100644
--- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMHADBTransactionBuffer.java
+++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMHADBTransactionBuffer.java
@@ -49,6 +49,8 @@ void flushIfNeeded(long snapshotWaitTime)
boolean shouldFlush(long snapshotWaitTime);
+ boolean flushIfPendingLimitReached() throws RocksDatabaseException, CodecException;
+
void init() throws RocksDatabaseException, CodecException;
void beginApplyingTransaction();
diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMHADBTransactionBufferImpl.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMHADBTransactionBufferImpl.java
index 07257ca09c38..0bda16fb7140 100644
--- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMHADBTransactionBufferImpl.java
+++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMHADBTransactionBufferImpl.java
@@ -17,6 +17,8 @@
package org.apache.hadoop.hdds.scm.ha;
+import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_HA_DBTRANSACTIONBUFFER_FLUSH_PENDING_LIMIT;
+import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_HA_DBTRANSACTIONBUFFER_FLUSH_PENDING_LIMIT_DEFAULT;
import static org.apache.hadoop.ozone.OzoneConsts.TRANSACTION_INFO_KEY;
import com.google.common.base.Preconditions;
@@ -24,6 +26,7 @@
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import java.util.concurrent.locks.ReentrantReadWriteLock;
+import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.hdds.scm.block.DeletedBlockLog;
import org.apache.hadoop.hdds.scm.block.DeletedBlockLogImpl;
import org.apache.hadoop.hdds.scm.metadata.SCMMetadataStore;
@@ -57,12 +60,27 @@ public class SCMHADBTransactionBufferImpl implements SCMHADBTransactionBuffer {
private final AtomicInteger applyingTransactions = new AtomicInteger(0);
private long lastSnapshotTimeMs = 0;
private final ReentrantReadWriteLock rwLock = new ReentrantReadWriteLock();
+ private final long flushPendingLimit;
public SCMHADBTransactionBufferImpl(StorageContainerManager scm) throws RocksDatabaseException, CodecException {
this.scm = scm;
+ this.flushPendingLimit = getFlushPendingLimit(scm.getConfiguration());
init();
}
+ private static long getFlushPendingLimit(OzoneConfiguration conf) {
+ long limit = conf.getLong(
+ OZONE_SCM_HA_DBTRANSACTIONBUFFER_FLUSH_PENDING_LIMIT,
+ OZONE_SCM_HA_DBTRANSACTIONBUFFER_FLUSH_PENDING_LIMIT_DEFAULT);
+ if (limit <= 0) {
+ LOG.warn("Invalid {}={} config value, using default {}",
+ OZONE_SCM_HA_DBTRANSACTIONBUFFER_FLUSH_PENDING_LIMIT, limit,
+ OZONE_SCM_HA_DBTRANSACTIONBUFFER_FLUSH_PENDING_LIMIT_DEFAULT);
+ limit = OZONE_SCM_HA_DBTRANSACTIONBUFFER_FLUSH_PENDING_LIMIT_DEFAULT;
+ }
+ return limit;
+ }
+
private BatchOperation getCurrentBatchOperation() {
return currentBatchOperation;
}
@@ -150,6 +168,29 @@ public void flushIfNeeded(long snapshotWaitTime)
}
}
+ /**
+ * Flushes the buffer if the number of pending DB operations has reached {@link #flushPendingLimit}.
+ *
+ * @return true if a flush was performed
+ */
+ @Override
+ public boolean flushIfPendingLimitReached() throws RocksDatabaseException, CodecException {
+ if (txFlushPending.get() < flushPendingLimit) {
+ return false;
+ }
+ rwLock.writeLock().lock();
+ try {
+ if (txFlushPending.get() < flushPendingLimit) {
+ return false;
+ }
+ LOG.debug("txFlushPending={} reached max limit={}, flushing", txFlushPending.get(), flushPendingLimit);
+ flushUnderWriteLock();
+ return true;
+ } finally {
+ rwLock.writeLock().unlock();
+ }
+ }
+
private void flushUnderWriteLock()
throws RocksDatabaseException, CodecException {
// write latest trx info into trx table in the same batch
diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMHADBTransactionBufferStub.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMHADBTransactionBufferStub.java
index 03b74330f77e..b7f1d901876d 100644
--- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMHADBTransactionBufferStub.java
+++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMHADBTransactionBufferStub.java
@@ -106,6 +106,11 @@ public void flushIfNeeded(long snapshotWaitTime) throws RocksDatabaseException {
flush();
}
+ @Override
+ public boolean flushIfPendingLimitReached() {
+ return false;
+ }
+
@Override
public void beginApplyingTransaction() {
}
diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMStateMachine.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMStateMachine.java
index 5036769ed3a6..8ec8f5127684 100644
--- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMStateMachine.java
+++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/SCMStateMachine.java
@@ -181,6 +181,8 @@ public CompletableFuture applyTransaction(
transactionBuffer.updateLatestTrxInfo(TransactionInfo.valueOf(appliedTermIndex));
updateLastAppliedTermIndex(appliedTermIndex);
+ transactionBuffer.flushIfPendingLimitReached();
+
// A restarted follower may catch up by applying data-carrying entries
// here rather than through notifyTermIndexUpdated, so check for catch-up
// in both places. No-op once the datanode protocol server has started.
diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/ha/TestSCMHADBTransactionBufferImpl.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/ha/TestSCMHADBTransactionBufferImpl.java
new file mode 100644
index 000000000000..d490fea1efe9
--- /dev/null
+++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/ha/TestSCMHADBTransactionBufferImpl.java
@@ -0,0 +1,333 @@
+/*
+ * 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.hadoop.hdds.scm.ha;
+
+import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_HA_DBTRANSACTIONBUFFER_FLUSH_PENDING_LIMIT;
+import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_HA_DBTRANSACTIONBUFFER_FLUSH_PENDING_LIMIT_DEFAULT;
+import static org.apache.hadoop.ozone.OzoneConsts.TRANSACTION_INFO_KEY;
+import static org.apache.ozone.test.GenericTestUtils.waitFor;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.clearInvocations;
+import static org.mockito.Mockito.doNothing;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+import com.google.protobuf.ByteString;
+import java.io.File;
+import java.time.Clock;
+import java.util.ArrayList;
+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.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import org.apache.hadoop.hdds.conf.OzoneConfiguration;
+import org.apache.hadoop.hdds.scm.block.BlockManager;
+import org.apache.hadoop.hdds.scm.block.DeletedBlockLogImpl;
+import org.apache.hadoop.hdds.scm.metadata.SCMMetadataStore;
+import org.apache.hadoop.hdds.scm.metadata.SCMMetadataStoreImpl;
+import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
+import org.apache.hadoop.hdds.utils.TransactionInfo;
+import org.apache.hadoop.hdds.utils.db.Table;
+import org.apache.hadoop.ozone.container.common.SCMTestUtils;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+
+/**
+ * Tests for the count-based flush of {@link SCMHADBTransactionBufferImpl}
+ * ({@link SCMHADBTransactionBufferImpl#flushIfPendingLimitReached()}).
+ */
+public class TestSCMHADBTransactionBufferImpl {
+
+ private static final TransactionInfo TRX_INFO_T4 = TransactionInfo.valueOf(1, 4);
+ private static final TransactionInfo TRX_INFO_T5 = TransactionInfo.valueOf(1, 5);
+ private static final ByteString VALUE = ByteString.copyFromUtf8("value");
+
+ @TempDir
+ private File testDir;
+
+ private SCMMetadataStore metadataStore;
+ private Table serviceConfigTable;
+ private Table transactionInfoTable;
+ private DeletedBlockLogImpl deletedBlockLog;
+ private final List buffers = new ArrayList<>();
+
+ @BeforeEach
+ public void setup() throws Exception {
+ metadataStore = new SCMMetadataStoreImpl(SCMTestUtils.getConf(testDir));
+ serviceConfigTable = metadataStore.getStatefulServiceConfigTable();
+ transactionInfoTable = metadataStore.getTransactionInfoTable();
+ deletedBlockLog = mock(DeletedBlockLogImpl.class);
+ }
+
+ @AfterEach
+ public void cleanup() throws Exception {
+ for (SCMHADBTransactionBufferImpl buffer : buffers) {
+ buffer.close();
+ }
+ if (metadataStore != null) {
+ metadataStore.stop();
+ }
+ }
+
+ @Test
+ public void testNoFlushBelowLimit() throws Exception {
+ SCMHADBTransactionBufferImpl buffer = newBuffer(3);
+ buffer.updateLatestTrxInfo(TRX_INFO_T4);
+ buffer.flush();
+ clearInvocations(deletedBlockLog);
+
+ for (int i = 0; i < 2; i++) {
+ put(buffer, "key" + i);
+ assertFalse(buffer.flushIfPendingLimitReached());
+ }
+ assertNotDurable("key0");
+ assertNotDurable("key1");
+ verify(deletedBlockLog, never()).onFlush();
+ }
+
+ @Test
+ public void testFlushAtLimitPersistsWholeBatchAndTransactionInfo() throws Exception {
+ SCMHADBTransactionBufferImpl buffer = newBuffer(3);
+ buffer.updateLatestTrxInfo(TRX_INFO_T4);
+ buffer.flush();
+ clearInvocations(deletedBlockLog);
+
+ put(buffer, "key0");
+ put(buffer, "key1");
+ buffer.updateLatestTrxInfo(TRX_INFO_T5);
+ put(buffer, "key2");
+
+ assertTrue(buffer.flushIfPendingLimitReached());
+ assertDurable("key0");
+ assertDurable("key1");
+ assertDurable("key2");
+ assertEquals(TRX_INFO_T5, transactionInfoTable.get(TRANSACTION_INFO_KEY));
+ verify(deletedBlockLog, times(1)).onFlush();
+ }
+
+ @Test
+ public void testSingleApplyOvershootingLimitFlushesOnce() throws Exception {
+ SCMHADBTransactionBufferImpl buffer = newBuffer(3);
+ buffer.updateLatestTrxInfo(TRX_INFO_T4);
+ buffer.flush();
+ clearInvocations(deletedBlockLog);
+
+ for (int i = 0; i < 5; i++) {
+ put(buffer, "key" + i);
+ }
+ buffer.updateLatestTrxInfo(TRX_INFO_T5);
+
+ assertTrue(buffer.flushIfPendingLimitReached());
+ for (int i = 0; i < 5; i++) {
+ assertDurable("key" + i);
+ }
+ assertEquals(TRX_INFO_T5, transactionInfoTable.get(TRANSACTION_INFO_KEY));
+ assertFalse(buffer.flushIfPendingLimitReached(), "buffer was drained, nothing left to flush");
+ verify(deletedBlockLog, times(1)).onFlush();
+ }
+
+ @Test
+ public void testCounterRestartsAfterFlush() throws Exception {
+ SCMHADBTransactionBufferImpl buffer = newBuffer(3);
+ buffer.updateLatestTrxInfo(TRX_INFO_T4);
+ buffer.flush();
+
+ for (int i = 0; i < 3; i++) {
+ put(buffer, "a" + i);
+ }
+ assertTrue(buffer.flushIfPendingLimitReached());
+
+ put(buffer, "b0");
+ put(buffer, "b1");
+ assertFalse(buffer.flushIfPendingLimitReached(), "only 2 pending after the previous flush");
+ assertNotDurable("b0");
+
+ put(buffer, "b2");
+ assertTrue(buffer.flushIfPendingLimitReached());
+ assertDurable("b0");
+ assertDurable("b1");
+ assertDurable("b2");
+ }
+
+ @Test
+ public void testFlushesWhileApplyingTransactionIsInProgress() throws Exception {
+ SCMHADBTransactionBufferImpl buffer = newBuffer(1);
+ buffer.updateLatestTrxInfo(TRX_INFO_T4);
+ buffer.flush();
+
+ buffer.beginApplyingTransaction();
+ try {
+ put(buffer, "key");
+ buffer.updateLatestTrxInfo(TRX_INFO_T5);
+ assertTrue(buffer.flushIfPendingLimitReached());
+ assertDurable("key");
+ assertEquals(TRX_INFO_T5, transactionInfoTable.get(TRANSACTION_INFO_KEY));
+ } finally {
+ buffer.endApplyingTransaction();
+ }
+ }
+
+ @ParameterizedTest
+ @ValueSource(longs = {0, -1, Long.MIN_VALUE})
+ public void testNonPositiveLimitFallsBackToDefault(long limit) throws Exception {
+ SCMHADBTransactionBufferImpl buffer = newBuffer(limit);
+ buffer.updateLatestTrxInfo(TRX_INFO_T4);
+ buffer.flush();
+
+ final long defaultLimit = OZONE_SCM_HA_DBTRANSACTIONBUFFER_FLUSH_PENDING_LIMIT_DEFAULT;
+ for (int i = 0; i < defaultLimit - 1; i++) {
+ put(buffer, "key" + i);
+ }
+ assertFalse(buffer.flushIfPendingLimitReached());
+ put(buffer, "key" + (defaultLimit - 1));
+ assertTrue(buffer.flushIfPendingLimitReached());
+ }
+
+ @Test
+ public void testConcurrentCallersFlushExactlyOnce() throws Exception {
+ SCMHADBTransactionBufferImpl buffer = newBuffer(1);
+ buffer.updateLatestTrxInfo(TRX_INFO_T4);
+ buffer.flush();
+ clearInvocations(deletedBlockLog);
+ put(buffer, "key");
+ buffer.updateLatestTrxInfo(TRX_INFO_T5);
+
+ final int threads = 8;
+ ExecutorService pool = Executors.newFixedThreadPool(threads);
+ try {
+ CountDownLatch start = new CountDownLatch(1);
+ AtomicInteger flushed = new AtomicInteger();
+ List> futures = new ArrayList<>();
+ for (int i = 0; i < threads; i++) {
+ futures.add(pool.submit(() -> {
+ start.await();
+ if (buffer.flushIfPendingLimitReached()) {
+ flushed.incrementAndGet();
+ }
+ return null;
+ }));
+ }
+ start.countDown();
+ for (Future> f : futures) {
+ f.get(10, TimeUnit.SECONDS);
+ }
+ assertEquals(1, flushed.get());
+ verify(deletedBlockLog, times(1)).onFlush();
+ assertDurable("key");
+ } finally {
+ pool.shutdownNow();
+ }
+ }
+
+ @Test
+ public void testFlushWaitsForBufferLock() throws Exception {
+ SCMHADBTransactionBufferImpl buffer = newBuffer(1);
+ buffer.updateLatestTrxInfo(TRX_INFO_T4);
+ buffer.flush();
+ put(buffer, "key");
+ buffer.updateLatestTrxInfo(TRX_INFO_T5);
+
+ AtomicInteger flushed = new AtomicInteger();
+ Thread flusher = new Thread(() -> {
+ try {
+ if (buffer.flushIfPendingLimitReached()) {
+ flushed.incrementAndGet();
+ }
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ });
+ buffer.lock();
+ try {
+ flusher.start();
+ waitFor(() -> flusher.getState() == Thread.State.WAITING, 10, 5_000);
+ assertNotDurable("key");
+ assertEquals(TRX_INFO_T4, transactionInfoTable.get(TRANSACTION_INFO_KEY));
+ } finally {
+ buffer.unlock();
+ }
+ flusher.join(10_000);
+ assertEquals(1, flushed.get());
+ assertDurable("key");
+ assertEquals(TRX_INFO_T5, transactionInfoTable.get(TRANSACTION_INFO_KEY));
+ }
+
+ @Test
+ public void testFailedFlushPropagatesAndIsRetried() throws Exception {
+ SCMHADBTransactionBufferImpl buffer = newBuffer(1);
+ buffer.updateLatestTrxInfo(TRX_INFO_T4);
+ buffer.flush();
+ put(buffer, "key");
+ buffer.updateLatestTrxInfo(TRX_INFO_T5);
+
+ doThrow(new IllegalStateException("injected")).when(deletedBlockLog).onFlush();
+ assertThrows(IllegalStateException.class, buffer::flushIfPendingLimitReached);
+
+ // Pending count was not reset by the failed flush, and the lock was released.
+ doNothing().when(deletedBlockLog).onFlush();
+ assertTrue(buffer.flushIfPendingLimitReached());
+ assertDurable("key");
+ assertEquals(TRX_INFO_T5, transactionInfoTable.get(TRANSACTION_INFO_KEY));
+ }
+
+ private SCMHADBTransactionBufferImpl newBuffer(OzoneConfiguration conf) throws Exception {
+ StorageContainerManager scm = mock(StorageContainerManager.class);
+ BlockManager blockManager = mock(BlockManager.class);
+ when(scm.getConfiguration()).thenReturn(conf);
+ when(scm.getScmMetadataStore()).thenReturn(metadataStore);
+ when(scm.getSystemClock()).thenReturn(Clock.systemUTC());
+ when(scm.getScmBlockManager()).thenReturn(blockManager);
+ when(blockManager.getDeletedBlockLog()).thenReturn(deletedBlockLog);
+ SCMHADBTransactionBufferImpl buffer = new SCMHADBTransactionBufferImpl(scm);
+ buffers.add(buffer);
+ return buffer;
+ }
+
+ private SCMHADBTransactionBufferImpl newBuffer(long limit) throws Exception {
+ OzoneConfiguration conf = new OzoneConfiguration();
+ conf.setLong(OZONE_SCM_HA_DBTRANSACTIONBUFFER_FLUSH_PENDING_LIMIT, limit);
+ return newBuffer(conf);
+ }
+
+ private void put(SCMHADBTransactionBufferImpl buffer, String key) throws Exception {
+ buffer.addToBuffer(serviceConfigTable, key, VALUE);
+ }
+
+ private void assertDurable(String key) throws Exception {
+ assertEquals(VALUE, serviceConfigTable.get(key), key + " should be durable");
+ }
+
+ private void assertNotDurable(String key) throws Exception {
+ assertNull(serviceConfigTable.get(key), key + " should still be buffered only");
+ }
+}
diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/ha/TestSCMHATransactionBufferMonitorTask.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/ha/TestSCMHATransactionBufferMonitorTask.java
index 9e0e971f1ca6..9a62601063eb 100644
--- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/ha/TestSCMHATransactionBufferMonitorTask.java
+++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/ha/TestSCMHATransactionBufferMonitorTask.java
@@ -94,6 +94,7 @@ public void setup() throws Exception {
deletedBlockLog = mock(DeletedBlockLogImpl.class);
Clock clock = mock(Clock.class);
when(clock.millis()).thenAnswer(invocation -> clockMillis.get());
+ when(scm.getConfiguration()).thenReturn(conf);
when(scm.getScmMetadataStore()).thenReturn(metadataStore);
when(scm.getSystemClock()).thenReturn(clock);
when(scm.getScmBlockManager()).thenReturn(blockManager);
diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/ha/TestSCMStateMachine.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/ha/TestSCMStateMachine.java
index 828606f2d424..0f6703cdf573 100644
--- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/ha/TestSCMStateMachine.java
+++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/ha/TestSCMStateMachine.java
@@ -17,16 +17,34 @@
package org.apache.hadoop.hdds.scm.ha;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.inOrder;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
+import java.util.concurrent.CompletableFuture;
+import org.apache.hadoop.hdds.protocol.proto.SCMRatisProtocol.RequestType;
import org.apache.hadoop.hdds.scm.container.placement.metrics.SCMMetrics;
+import org.apache.hadoop.hdds.scm.exceptions.SCMException;
+import org.apache.hadoop.hdds.scm.exceptions.SCMException.ResultCodes;
+import org.apache.hadoop.hdds.scm.ha.invoker.ScmInvoker;
+import org.apache.hadoop.hdds.scm.safemode.SCMSafeModeManager;
+import org.apache.hadoop.hdds.scm.server.SCMDatanodeProtocolServer;
import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
import org.apache.hadoop.hdds.utils.TransactionInfo;
import org.apache.ratis.proto.RaftProtos;
+import org.apache.ratis.protocol.Message;
import org.apache.ratis.server.protocol.TermIndex;
+import org.apache.ratis.statemachine.TransactionContext;
+import org.apache.ratis.util.ExitUtils;
import org.junit.jupiter.api.Test;
+import org.mockito.InOrder;
/**
* Test SCMStateMachine events recording.
@@ -49,4 +67,101 @@ public void testRatisEventsRecording() throws Exception {
metrics.unRegister();
}
+
+ @Test
+ public void testApplyTransactionFlushesAfterRecordingTransactionInfo() throws Exception {
+ SCMHADBTransactionBuffer buffer = mock(SCMHADBTransactionBuffer.class);
+ SCMStateMachine stateMachine = newStateMachine(buffer, succeedingInvoker());
+
+ CompletableFuture result = stateMachine.applyTransaction(newTransaction(7));
+
+ assertTrue(result.isDone() && !result.isCompletedExceptionally());
+ InOrder order = inOrder(buffer);
+ order.verify(buffer).beginApplyingTransaction();
+ order.verify(buffer).updateLatestTrxInfo(TransactionInfo.valueOf(TermIndex.valueOf(1, 7)));
+ order.verify(buffer).flushIfPendingLimitReached();
+ order.verify(buffer).endApplyingTransaction();
+ }
+
+ @Test
+ public void testFlushCheckedForEveryAppliedTransaction() throws Exception {
+ SCMHADBTransactionBuffer buffer = mock(SCMHADBTransactionBuffer.class);
+ SCMStateMachine stateMachine = newStateMachine(buffer, succeedingInvoker());
+
+ for (int i = 1; i <= 5; i++) {
+ stateMachine.applyTransaction(newTransaction(i));
+ }
+
+ verify(buffer, times(5)).flushIfPendingLimitReached();
+ verify(buffer, times(5)).endApplyingTransaction();
+ }
+
+ @Test
+ public void testFlushStillCheckedWhenTransactionIsLogicallyRejected() throws Exception {
+ SCMHADBTransactionBuffer buffer = mock(SCMHADBTransactionBuffer.class);
+ ScmInvoker> invoker = mock(ScmInvoker.class);
+ when(invoker.invokeLocal(anyString(), any()))
+ .thenThrow(new SCMException("rejected", ResultCodes.FAILED_TO_FIND_CONTAINER));
+ SCMStateMachine stateMachine = newStateMachine(buffer, invoker);
+
+ CompletableFuture result = stateMachine.applyTransaction(newTransaction(3));
+
+ assertTrue(result.isCompletedExceptionally());
+ InOrder order = inOrder(buffer);
+ order.verify(buffer).updateLatestTrxInfo(TransactionInfo.valueOf(TermIndex.valueOf(1, 3)));
+ order.verify(buffer).flushIfPendingLimitReached();
+ order.verify(buffer).endApplyingTransaction();
+ }
+
+ @Test
+ public void testFlushFailureTerminatesAndClosesApplyingWindow() throws Exception {
+ ExitUtils.disableSystemExit();
+ try {
+ SCMHADBTransactionBuffer buffer = mock(SCMHADBTransactionBuffer.class);
+ doThrow(new IllegalStateException("injected flush failure")).when(buffer).flushIfPendingLimitReached();
+ SCMStateMachine stateMachine = newStateMachine(buffer, succeedingInvoker());
+
+ assertThrows(ExitUtils.ExitException.class, () -> stateMachine.applyTransaction(newTransaction(9)));
+
+ verify(buffer).endApplyingTransaction();
+ } finally {
+ ExitUtils.clear();
+ }
+ }
+
+ private static TransactionContext newTransaction(long index) throws Exception {
+ SCMRatisRequest request = SCMRatisRequest.of(RequestType.PIPELINE, "op", new Class>[0]);
+ RaftProtos.StateMachineLogEntryProto smLogEntry = RaftProtos.StateMachineLogEntryProto.newBuilder()
+ .setLogData(request.encode().getContent())
+ .build();
+ RaftProtos.LogEntryProto logEntry = RaftProtos.LogEntryProto.newBuilder()
+ .setTerm(1)
+ .setIndex(index)
+ .setStateMachineLogEntry(smLogEntry)
+ .build();
+ TransactionContext trx = mock(TransactionContext.class);
+ when(trx.getStateMachineLogEntry()).thenReturn(smLogEntry);
+ when(trx.getLogEntry()).thenReturn(logEntry);
+ return trx;
+ }
+
+ private static SCMStateMachine newStateMachine(SCMHADBTransactionBuffer buffer, ScmInvoker> invoker) {
+ StorageContainerManager scm = mock(StorageContainerManager.class);
+ when(scm.getMetrics()).thenReturn(mock(SCMMetrics.class));
+ SCMContext context = mock(SCMContext.class);
+ when(context.isLeader()).thenReturn(true);
+ when(scm.getScmContext()).thenReturn(context);
+ when(scm.getDatanodeProtocolServer()).thenReturn(mock(SCMDatanodeProtocolServer.class));
+ when(scm.getScmSafeModeManager()).thenReturn(mock(SCMSafeModeManager.class));
+ when(buffer.getLatestTrxInfo()).thenReturn(TransactionInfo.valueOf(TermIndex.valueOf(0, 0)));
+ SCMStateMachine stateMachine = new SCMStateMachine(scm, buffer);
+ stateMachine.registerInvoker(RequestType.PIPELINE, invoker);
+ return stateMachine;
+ }
+
+ private static ScmInvoker> succeedingInvoker() throws Exception {
+ ScmInvoker> invoker = mock(ScmInvoker.class);
+ when(invoker.invokeLocal(anyString(), any())).thenReturn(Message.EMPTY);
+ return invoker;
+ }
}