Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
11 changes: 10 additions & 1 deletion hadoop-hdds/common/src/main/resources/ozone-default.xml
Original file line number Diff line number Diff line change
Expand Up @@ -4389,7 +4389,16 @@
<tag>SCM, OZONE</tag>
<description>Wait duration for flush of buffered transaction.</description>
</property>

<property>
<name>ozone.scm.ha.dbtransactionbuffer.flush.pending.limit</name>
<value>10000</value>
<tag>SCM, OZONE</tag>
<description>
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.
</description>
</property>

<property>
<name>ozone.s3g.kerberos.keytab.file</name>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,8 @@ void flushIfNeeded(long snapshotWaitTime)

boolean shouldFlush(long snapshotWaitTime);

boolean flushIfPendingLimitReached() throws RocksDatabaseException, CodecException;

void init() throws RocksDatabaseException, CodecException;

void beginApplyingTransaction();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,13 +17,16 @@

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;
import java.util.concurrent.atomic.AtomicInteger;
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;
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,11 @@ public void flushIfNeeded(long snapshotWaitTime) throws RocksDatabaseException {
flush();
}

@Override
public boolean flushIfPendingLimitReached() {
return false;
}

@Override
public void beginApplyingTransaction() {
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -181,6 +181,8 @@ public CompletableFuture<Message> 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.
Expand Down
Loading
Loading