From 8e23577e0324b1ef3105c7f707c0544481191f39 Mon Sep 17 00:00:00 2001 From: sarvekshayr Date: Mon, 5 Oct 2026 14:14:50 +0530 Subject: [PATCH 1/3] HDDS-16666. Flush SCM transaction in memory during apply transaction --- .../apache/hadoop/hdds/scm/ScmConfigKeys.java | 4 + .../src/main/resources/ozone-default.xml | 11 +- .../hdds/scm/ha/SCMHADBTransactionBuffer.java | 2 + .../scm/ha/SCMHADBTransactionBufferImpl.java | 41 +++ .../scm/ha/SCMHADBTransactionBufferStub.java | 5 + .../hadoop/hdds/scm/ha/SCMStateMachine.java | 2 + .../ha/TestSCMHADBTransactionBufferImpl.java | 333 ++++++++++++++++++ ...TestSCMHATransactionBufferMonitorTask.java | 1 + .../hdds/scm/ha/TestSCMStateMachine.java | 115 ++++++ 9 files changed, 513 insertions(+), 1 deletion(-) create mode 100644 hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/ha/TestSCMHADBTransactionBufferImpl.java 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..6b0e706ae98b 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 (flushPendingLimit <= 0 || 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..c7d82b8ab822 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(SCMMetrics.create()); + 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; + } } From 5e42dcb6ed691d3e7495cb9d28c79aaab8091ea2 Mon Sep 17 00:00:00 2001 From: sarvekshayr Date: Mon, 5 Oct 2026 15:34:13 +0530 Subject: [PATCH 2/3] Redundant check removed --- .../apache/hadoop/hdds/scm/ha/SCMHADBTransactionBufferImpl.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 6b0e706ae98b..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 @@ -175,7 +175,7 @@ public void flushIfNeeded(long snapshotWaitTime) */ @Override public boolean flushIfPendingLimitReached() throws RocksDatabaseException, CodecException { - if (flushPendingLimit <= 0 || txFlushPending.get() < flushPendingLimit) { + if (txFlushPending.get() < flushPendingLimit) { return false; } rwLock.writeLock().lock(); From 108e329c1b63da4a3dd03fe42223b4420f04027a Mon Sep 17 00:00:00 2001 From: Sarveksha Yeshavantha Raju <79865743+sarvekshayr@users.noreply.github.com> Date: Mon, 5 Oct 2026 18:53:08 +0530 Subject: [PATCH 3/3] Prevent duplicate SCMMetrics source registration in tests Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- .../java/org/apache/hadoop/hdds/scm/ha/TestSCMStateMachine.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 c7d82b8ab822..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 @@ -147,7 +147,7 @@ private static TransactionContext newTransaction(long index) throws Exception { private static SCMStateMachine newStateMachine(SCMHADBTransactionBuffer buffer, ScmInvoker invoker) { StorageContainerManager scm = mock(StorageContainerManager.class); - when(scm.getMetrics()).thenReturn(SCMMetrics.create()); + when(scm.getMetrics()).thenReturn(mock(SCMMetrics.class)); SCMContext context = mock(SCMContext.class); when(context.isLeader()).thenReturn(true); when(scm.getScmContext()).thenReturn(context);