rich7420 commented on code in PR #11118:
URL: https://github.com/apache/ozone/pull/11118#discussion_r3861474305
##########
hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMDirectoriesPurgeRequestWithFSO.java:
##########
@@ -134,93 +135,17 @@ public OMClientResponse
validateAndUpdateCache(OzoneManager ozoneManager, Execut
AUDIT.logWriteFailure(ozoneManager.buildAuditMessageForFailure(OMSystemAction.DIRECTORY_DELETION,
null, oe));
return new
OMDirectoriesPurgeResponseWithFSO(createErrorOMResponse(omResponse, oe));
}
+ PurgeApplyResult result = new PurgeApplyResult();
try {
- int numSubDirMoved = 0, numSubFilesMoved = 0, numDirsDeleted = 0;
- Map<VolumeBucketId, BucketNameInfo> volumeBucketIdMap =
purgeDirsRequest.getBucketNameInfosList().stream()
- .collect(Collectors.toMap(bucketNameInfo ->
- new VolumeBucketId(bucketNameInfo.getVolumeId(),
bucketNameInfo.getBucketId()),
- Function.identity()));
- for (OzoneManagerProtocolProtos.PurgePathRequest path : purgeRequests) {
- for (OzoneManagerProtocolProtos.KeyInfo key :
path.getMarkDeletedSubDirsList()) {
- ProcessedKeyInfo processed = processDeleteKey(key, path,
omMetadataManager);
- subDirNames.add(processed.deleteKey);
-
- omMetrics.decNumKeys();
- omMetrics.incNumKeyDeletesInternal();
- OmBucketInfo omBucketInfo = getBucketInfo(omMetadataManager,
- processed.volumeName, processed.bucketName);
- // bucketInfo can be null in case of delete volume or bucket
- // or key does not belong to bucket as bucket is recreated
- if (null != omBucketInfo && omBucketInfo.getObjectID() ==
path.getBucketId()) {
- omBucketInfo.decrUsedNamespace(1L, true);
- String ozoneDbKey =
omMetadataManager.getOzonePathKey(path.getVolumeId(),
- path.getBucketId(), processed.keyInfo.getParentObjectID(),
- processed.keyInfo.getFileName());
- omMetadataManager.getDirectoryTable().addCacheEntry(new
CacheKey<>(ozoneDbKey),
- CacheValue.get(context.getIndex()));
- volBucketInfoMap.putIfAbsent(processed.volBucketPair,
omBucketInfo);
- }
- }
-
- for (OzoneManagerProtocolProtos.KeyInfo key :
path.getDeletedSubFilesList()) {
- ProcessedKeyInfo processed = processDeleteKey(key, path,
omMetadataManager);
- subFileNames.add(processed.deleteKey);
-
- // If omKeyInfo has hsync metadata, delete its corresponding open
key as well
- String dbOpenKey;
- String hsyncClientId =
processed.keyInfo.getMetadata().get(OzoneConsts.HSYNC_CLIENT_ID);
- if (hsyncClientId != null) {
- long parentId = processed.keyInfo.getParentObjectID();
- dbOpenKey = omMetadataManager.getOpenFileName(path.getVolumeId(),
path.getBucketId(),
- parentId, processed.keyInfo.getFileName(), hsyncClientId);
- OmKeyInfo openKeyInfo =
omMetadataManager.getOpenKeyTable(getBucketLayout()).get(dbOpenKey);
- if (openKeyInfo != null) {
- openKeyInfo = openKeyInfo.withMetadataMutations(
- metadata -> metadata.put(DELETED_HSYNC_KEY, "true"));
- openKeyInfoMap.put(dbOpenKey, openKeyInfo);
- }
- }
-
- omMetrics.decNumKeys();
- omMetrics.incNumKeyDeletesInternal();
- numSubFilesMoved++;
- OmBucketInfo omBucketInfo = getBucketInfo(omMetadataManager,
- processed.volumeName, processed.bucketName);
- // bucketInfo can be null in case of delete volume or bucket
- // or key does not belong to bucket as bucket is recreated
- if (null != omBucketInfo
- && omBucketInfo.getObjectID() == path.getBucketId()) {
- long totalSize = sumBlockLengths(processed.keyInfo);
- omBucketInfo.decrUsedBytes(totalSize, true);
- omBucketInfo.decrUsedNamespace(1L, true);
- String ozoneDbKey =
omMetadataManager.getOzonePathKey(path.getVolumeId(),
- path.getBucketId(), processed.keyInfo.getParentObjectID(),
- processed.keyInfo.getFileName());
- omMetadataManager.getFileTable().addCacheEntry(new
CacheKey<>(ozoneDbKey),
- CacheValue.get(context.getIndex()));
- volBucketInfoMap.putIfAbsent(processed.volBucketPair,
omBucketInfo);
- }
- }
- if (path.hasDeletedDir()) {
- deletedDirNames.add(path.getDeletedDir());
- BucketNameInfo bucketNameInfo = volumeBucketIdMap.get(new
VolumeBucketId(path.getVolumeId(),
- path.getBucketId()));
- OmBucketInfo omBucketInfo = getBucketInfo(omMetadataManager,
- bucketNameInfo.getVolumeName(), bucketNameInfo.getBucketName());
- if (omBucketInfo != null && omBucketInfo.getObjectID() ==
path.getBucketId()) {
- omBucketInfo.purgeSnapshotUsedNamespace(1);
- volBucketInfoMap.put(Pair.of(omBucketInfo.getVolumeName(),
omBucketInfo.getBucketName()), omBucketInfo);
- }
- numDirsDeleted++;
- }
- }
+ // Phase 2 (under the bucket write lock): apply the prepared cache
tombstones, hsync open-key cleanup and
+ // aggregated per-bucket quota changes.
+ applyPreparedEntries(preparedSubDirs, preparedSubFiles,
preparedDirPurges, omMetadataManager,
Review Comment:
Could we split a dense `PurgePathRequest` before applying it under the
bucket lock? The apply path still processes every prepared entry under one lock
acquisition, while DDS batching keeps each path request indivisible. This
leaves the single-dense-directory case from HDDS-16220 unresolved.
##########
hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/service/TestOmMixedWorkloadUnderDeletionBench.java:
##########
@@ -0,0 +1,686 @@
+/*
+ * 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.ozone.om.service;
+
+import static
org.apache.hadoop.fs.CommonConfigurationKeysPublic.FS_DEFAULT_NAME_KEY;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_BLOCK_DELETING_SERVICE_INTERVAL;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_FS_ITERATE_BATCH_SIZE;
+
+import java.io.File;
+import java.io.IOException;
+import java.lang.reflect.Method;
+import java.net.URL;
+import java.net.URLClassLoader;
+import java.security.AccessController;
+import java.security.PrivilegedActionException;
+import java.security.PrivilegedExceptionAction;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Locale;
+import java.util.concurrent.Callable;
+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.AtomicBoolean;
+import org.apache.hadoop.fs.FSDataInputStream;
+import org.apache.hadoop.fs.FSDataOutputStream;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.hdds.conf.OzoneConfiguration;
+import org.apache.hadoop.hdds.protocol.StorageType;
+import org.apache.hadoop.hdds.utils.db.CodecBuffer;
+import org.apache.hadoop.hdds.utils.db.Table;
+import org.apache.hadoop.ozone.MiniOzoneCluster;
+import org.apache.hadoop.ozone.OzoneConsts;
+import org.apache.hadoop.ozone.client.BucketArgs;
+import org.apache.hadoop.ozone.client.OzoneBucket;
+import org.apache.hadoop.ozone.client.OzoneClient;
+import org.apache.hadoop.ozone.client.OzoneVolume;
+import org.apache.hadoop.ozone.om.OMConfigKeys;
+import org.apache.hadoop.ozone.om.OMMetadataManager;
+import org.apache.hadoop.ozone.om.helpers.BucketLayout;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * End-to-end benchmark reproducing the interactive-workload degradation seen
when a large background deletion backlog
+ * is being reclaimed on the same bucket that clients are actively writing to.
+ *
+ * <p>Every OM write is applied on a single serial state-machine apply thread
under a per-bucket write lock. Reclaiming
+ * a deletion backlog drives {@code PurgeDirectories} transactions through
that thread; when the deleted directories are
+ * densely populated, a single purge batch moves a large number of
sub-files/sub-dirs under one write-lock hold, so the
+ * apply thread — and the bucket write lock — is occupied for the whole batch.
Concurrent user {@code create},
+ * {@code mkdir} and {@code rename} contend for that same lock and thread, and
reads contend for the bucket read lock,
+ * so their latency degrades while the backlog drains.
+ *
+ * <p>This benchmark stages two sets under a single volume and FSO bucket: a
densely-populated <em>backlog</em> subtree
+ * that is recursively deleted and drained, and a separate stable
<em>workload</em> set that is never deleted during the
+ * run (a pre-staged dataset the read ops resolve against, plus a per-thread
scratch area the write ops create into).
+ * Sharing the bucket is intentional — deletion and the client workload
contend on the same bucket lock, while the
+ * workload always hits live paths. It then compares a mixed client workload
on that bucket — create, mkdir, rename and
+ * a data-plane file write (create + write a block), plus a data-plane file
read and the metadata read RPCs that take
+ * the bucket read lock (getFileStatus, listStatus, getBucketInfo, lookupKey)
— in two conditions:
+ * <ul>
+ * <li><b>control</b> — no deletion running, and</li>
+ * <li><b>under load</b> — the same workload while the backlog subtree is
recursively deleted and fully purged from
+ * OM in the background (both FSO phases: moved into the deletedTable,
then purged back out),</li>
+ * </ul>
+ * reporting per-operation p50/p99 latency and the under-load degradation, so
the two code versions can be compared on
+ * how much apply-thread purge work bleeds into interactive latency. It also
reports the phase-1 drain time — the
+ * {@code DirectoryDeletingService} move into the deletedTable via {@code
OMDirectoriesPurgeRequestWithFSO}, the apply
+ * path this change optimizes — separately from the full both-phase drain, so
pure apply-thread throughput can be
+ * compared alongside the interactive degradation.
+ *
+ * <p>Deletion is configured with production-representative per-task limits so
batches are large. The {@code benchmark}
+ * tag is excluded from {@code mvn test} and CI by default, so it must be
re-enabled explicitly to run on demand
+ * (rebuild the reactor first to avoid stale-class errors):
+ * <pre>
+ * mvn -pl :ozone-integration-test test -DskipShade -DskipRecon \
+ * -Dtest=TestOmMixedWorkloadUnderDeletionBench -Dgroups=benchmark
-Dexcluded-test-groups= \
+ * -Dsurefire.failIfNoSpecifiedTests=false
+ * </pre>
+ * Tunables: {@code bench.backlogDirs} (default 80), {@code
bench.backlogFilesPerDir} (default 1000),
+ * {@code bench.backlogNonEmptyEvery} (default 3 — every 3rd backlog file is
written with a block, the rest are
+ * empty so a large backlog stays cheap to stage), {@code bench.workloadDirs}
(default 20) and
+ * {@code bench.workloadFilesPerDir} (default 100) sizing the stable dataset
the read ops resolve against,
+ * {@code bench.fileBytes} (default 1 MiB) sizing the data-plane file
write/read payload (and the block-bearing
+ * staged files the reads pull), {@code bench.clientThreads} (default 4),
+ * {@code bench.opsPerThread} (default 400), {@code
bench.pathDeletingLimitPerTask} (default 2000) and
+ * {@code bench.keyDeletingLimitPerTask} (default 40000). The last two size
how much a single deletion round gathers;
+ * with the Ratis appender byte limit non-binding at these entry sizes, a
round's paths pack into one purge
+ * transaction, so raising them makes each apply move far more entries under a
single bucket write-lock hold — the
+ * regime where the apply-thread per-entry cost dominates interactive latency.
+ *
+ * <p>Adding {@code -Dbench.profile.event=<cpu|lock|wall|alloc>} profiles only
the under-load window with
+ * async-profiler, loaded reflectively from a local install whose paths must
be supplied via
+ * {@code -Dbench.profiler.jar} (the async-profiler jar) and {@code
-Dbench.profiler.lib} (its native library); the
+ * JFR is written under {@code -Dbench.profile.out}, default {@code /tmp}. For
accurate leaf frames also pass
+ * {@code -DargLine="-XX:+UnlockDiagnosticVMOptions -XX:+DebugNonSafepoints"}.
+ */
+@Tag("benchmark")
+public class TestOmMixedWorkloadUnderDeletionBench {
+
+ private static final Logger LOG =
LoggerFactory.getLogger(TestOmMixedWorkloadUnderDeletionBench.class);
+
+ private static final String OP_CREATE = "create";
+ private static final String OP_MKDIR = "mkdir";
+ private static final String OP_RENAME = "rename";
+ private static final String OP_FILEWRITE = "filewrite";
+ private static final String OP_FILEREAD = "fileread";
+ private static final String OP_GETFILESTATUS = "getfilestatus";
+ private static final String OP_LISTSTATUS = "liststatus";
+ private static final String OP_INFOBUCKET = "infobucket";
+ private static final String OP_GETKEYINFO = "getkeyinfo";
+ private static final String[] OPS =
+ {OP_CREATE, OP_MKDIR, OP_RENAME, OP_FILEWRITE, OP_FILEREAD,
+ OP_GETFILESTATUS, OP_LISTSTATUS, OP_INFOBUCKET, OP_GETKEYINFO};
+
+ // The three sandboxes share the parent /workload/bucket but are separate
subtrees; only the backlog is deleted.
+ // The deletion set (built, recursively deleted, then drained through the
apply thread).
+ private static final String BACKLOG_ROOT = "workload/bucket/backlog";
+ // The workload set — never deleted during the test: a stable pre-staged
dataset the read ops resolve against, plus
+ // a scratch area the write ops create into. Keeping this separate from the
deletion set means the concurrent
+ // operations always hit live paths while the backlog drains.
+ private static final String WORKLOAD_DATA_ROOT = "workload/bucket/data";
+ private static final String WORKLOAD_SCRATCH_ROOT =
"workload/bucket/scratch";
+
+ // Every bench.backlogNonEmptyEvery-th backlog file is written with a single
block so its KeyInfo carries a
+ // key-location list, exercising the block-metadata parse/serialize the
purge apply and flush paths hit in
+ // production. The rest are left empty because block-bearing files are far
more expensive to stage (block
+ // allocation + datanode write + commit), and a large backlog is what
actually stresses the apply thread.
+ private static final byte[] FILE_CONTENT = new byte[4];
+
+ /**
+ * Removes test-harness-only overhead that would otherwise distort the
apply/flush cost under measurement: the
+ * mini-cluster unconditionally enables {@link CodecBuffer} leak detection
(a per-allocation finalizer), and the test
+ * log config runs the {@code CodecBuffer}/managed-RocksDB loggers at
DEBUG/TRACE, which capture a full stack trace on
+ * every buffer allocation. Neither happens in a production OM running at
INFO.
+ */
+ private static void stripTestOnlyOverhead() {
+ CodecBuffer.disableLeakDetection();
+
org.apache.log4j.Logger.getLogger("org.apache.hadoop.hdds.utils.db.CodecBuffer")
+ .setLevel(org.apache.log4j.Level.INFO);
+
org.apache.log4j.Logger.getLogger("org.apache.hadoop.hdds.utils.db.managed")
+ .setLevel(org.apache.log4j.Level.INFO);
+ }
+
+ @Test
+ @Timeout(value = 120, unit = TimeUnit.MINUTES)
Review Comment:
Could we drop this per-test `@Timeout` per HDDS-12575? The benchmark already
has an explicit phase-drain deadline.
##########
hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/service/TestOmMixedWorkloadUnderDeletionBench.java:
##########
@@ -0,0 +1,686 @@
+/*
+ * 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.ozone.om.service;
+
+import static
org.apache.hadoop.fs.CommonConfigurationKeysPublic.FS_DEFAULT_NAME_KEY;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_BLOCK_DELETING_SERVICE_INTERVAL;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_FS_ITERATE_BATCH_SIZE;
+
+import java.io.File;
+import java.io.IOException;
+import java.lang.reflect.Method;
+import java.net.URL;
+import java.net.URLClassLoader;
+import java.security.AccessController;
+import java.security.PrivilegedActionException;
+import java.security.PrivilegedExceptionAction;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Locale;
+import java.util.concurrent.Callable;
+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.AtomicBoolean;
+import org.apache.hadoop.fs.FSDataInputStream;
+import org.apache.hadoop.fs.FSDataOutputStream;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.hdds.conf.OzoneConfiguration;
+import org.apache.hadoop.hdds.protocol.StorageType;
+import org.apache.hadoop.hdds.utils.db.CodecBuffer;
+import org.apache.hadoop.hdds.utils.db.Table;
+import org.apache.hadoop.ozone.MiniOzoneCluster;
+import org.apache.hadoop.ozone.OzoneConsts;
+import org.apache.hadoop.ozone.client.BucketArgs;
+import org.apache.hadoop.ozone.client.OzoneBucket;
+import org.apache.hadoop.ozone.client.OzoneClient;
+import org.apache.hadoop.ozone.client.OzoneVolume;
+import org.apache.hadoop.ozone.om.OMConfigKeys;
+import org.apache.hadoop.ozone.om.OMMetadataManager;
+import org.apache.hadoop.ozone.om.helpers.BucketLayout;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * End-to-end benchmark reproducing the interactive-workload degradation seen
when a large background deletion backlog
+ * is being reclaimed on the same bucket that clients are actively writing to.
+ *
+ * <p>Every OM write is applied on a single serial state-machine apply thread
under a per-bucket write lock. Reclaiming
+ * a deletion backlog drives {@code PurgeDirectories} transactions through
that thread; when the deleted directories are
+ * densely populated, a single purge batch moves a large number of
sub-files/sub-dirs under one write-lock hold, so the
+ * apply thread — and the bucket write lock — is occupied for the whole batch.
Concurrent user {@code create},
+ * {@code mkdir} and {@code rename} contend for that same lock and thread, and
reads contend for the bucket read lock,
+ * so their latency degrades while the backlog drains.
+ *
+ * <p>This benchmark stages two sets under a single volume and FSO bucket: a
densely-populated <em>backlog</em> subtree
+ * that is recursively deleted and drained, and a separate stable
<em>workload</em> set that is never deleted during the
+ * run (a pre-staged dataset the read ops resolve against, plus a per-thread
scratch area the write ops create into).
+ * Sharing the bucket is intentional — deletion and the client workload
contend on the same bucket lock, while the
+ * workload always hits live paths. It then compares a mixed client workload
on that bucket — create, mkdir, rename and
+ * a data-plane file write (create + write a block), plus a data-plane file
read and the metadata read RPCs that take
+ * the bucket read lock (getFileStatus, listStatus, getBucketInfo, lookupKey)
— in two conditions:
+ * <ul>
+ * <li><b>control</b> — no deletion running, and</li>
+ * <li><b>under load</b> — the same workload while the backlog subtree is
recursively deleted and fully purged from
+ * OM in the background (both FSO phases: moved into the deletedTable,
then purged back out),</li>
+ * </ul>
+ * reporting per-operation p50/p99 latency and the under-load degradation, so
the two code versions can be compared on
+ * how much apply-thread purge work bleeds into interactive latency. It also
reports the phase-1 drain time — the
+ * {@code DirectoryDeletingService} move into the deletedTable via {@code
OMDirectoriesPurgeRequestWithFSO}, the apply
+ * path this change optimizes — separately from the full both-phase drain, so
pure apply-thread throughput can be
+ * compared alongside the interactive degradation.
+ *
+ * <p>Deletion is configured with production-representative per-task limits so
batches are large. The {@code benchmark}
+ * tag is excluded from {@code mvn test} and CI by default, so it must be
re-enabled explicitly to run on demand
+ * (rebuild the reactor first to avoid stale-class errors):
+ * <pre>
+ * mvn -pl :ozone-integration-test test -DskipShade -DskipRecon \
+ * -Dtest=TestOmMixedWorkloadUnderDeletionBench -Dgroups=benchmark
-Dexcluded-test-groups= \
+ * -Dsurefire.failIfNoSpecifiedTests=false
+ * </pre>
+ * Tunables: {@code bench.backlogDirs} (default 80), {@code
bench.backlogFilesPerDir} (default 1000),
+ * {@code bench.backlogNonEmptyEvery} (default 3 — every 3rd backlog file is
written with a block, the rest are
+ * empty so a large backlog stays cheap to stage), {@code bench.workloadDirs}
(default 20) and
+ * {@code bench.workloadFilesPerDir} (default 100) sizing the stable dataset
the read ops resolve against,
+ * {@code bench.fileBytes} (default 1 MiB) sizing the data-plane file
write/read payload (and the block-bearing
+ * staged files the reads pull), {@code bench.clientThreads} (default 4),
+ * {@code bench.opsPerThread} (default 400), {@code
bench.pathDeletingLimitPerTask} (default 2000) and
+ * {@code bench.keyDeletingLimitPerTask} (default 40000). The last two size
how much a single deletion round gathers;
+ * with the Ratis appender byte limit non-binding at these entry sizes, a
round's paths pack into one purge
+ * transaction, so raising them makes each apply move far more entries under a
single bucket write-lock hold — the
+ * regime where the apply-thread per-entry cost dominates interactive latency.
+ *
+ * <p>Adding {@code -Dbench.profile.event=<cpu|lock|wall|alloc>} profiles only
the under-load window with
+ * async-profiler, loaded reflectively from a local install whose paths must
be supplied via
+ * {@code -Dbench.profiler.jar} (the async-profiler jar) and {@code
-Dbench.profiler.lib} (its native library); the
+ * JFR is written under {@code -Dbench.profile.out}, default {@code /tmp}. For
accurate leaf frames also pass
+ * {@code -DargLine="-XX:+UnlockDiagnosticVMOptions -XX:+DebugNonSafepoints"}.
+ */
+@Tag("benchmark")
+public class TestOmMixedWorkloadUnderDeletionBench {
+
+ private static final Logger LOG =
LoggerFactory.getLogger(TestOmMixedWorkloadUnderDeletionBench.class);
+
+ private static final String OP_CREATE = "create";
+ private static final String OP_MKDIR = "mkdir";
+ private static final String OP_RENAME = "rename";
+ private static final String OP_FILEWRITE = "filewrite";
+ private static final String OP_FILEREAD = "fileread";
+ private static final String OP_GETFILESTATUS = "getfilestatus";
+ private static final String OP_LISTSTATUS = "liststatus";
+ private static final String OP_INFOBUCKET = "infobucket";
+ private static final String OP_GETKEYINFO = "getkeyinfo";
+ private static final String[] OPS =
+ {OP_CREATE, OP_MKDIR, OP_RENAME, OP_FILEWRITE, OP_FILEREAD,
+ OP_GETFILESTATUS, OP_LISTSTATUS, OP_INFOBUCKET, OP_GETKEYINFO};
+
+ // The three sandboxes share the parent /workload/bucket but are separate
subtrees; only the backlog is deleted.
+ // The deletion set (built, recursively deleted, then drained through the
apply thread).
+ private static final String BACKLOG_ROOT = "workload/bucket/backlog";
+ // The workload set — never deleted during the test: a stable pre-staged
dataset the read ops resolve against, plus
+ // a scratch area the write ops create into. Keeping this separate from the
deletion set means the concurrent
+ // operations always hit live paths while the backlog drains.
+ private static final String WORKLOAD_DATA_ROOT = "workload/bucket/data";
+ private static final String WORKLOAD_SCRATCH_ROOT =
"workload/bucket/scratch";
+
+ // Every bench.backlogNonEmptyEvery-th backlog file is written with a single
block so its KeyInfo carries a
+ // key-location list, exercising the block-metadata parse/serialize the
purge apply and flush paths hit in
+ // production. The rest are left empty because block-bearing files are far
more expensive to stage (block
+ // allocation + datanode write + commit), and a large backlog is what
actually stresses the apply thread.
+ private static final byte[] FILE_CONTENT = new byte[4];
+
+ /**
+ * Removes test-harness-only overhead that would otherwise distort the
apply/flush cost under measurement: the
+ * mini-cluster unconditionally enables {@link CodecBuffer} leak detection
(a per-allocation finalizer), and the test
+ * log config runs the {@code CodecBuffer}/managed-RocksDB loggers at
DEBUG/TRACE, which capture a full stack trace on
+ * every buffer allocation. Neither happens in a production OM running at
INFO.
+ */
+ private static void stripTestOnlyOverhead() {
+ CodecBuffer.disableLeakDetection();
+
org.apache.log4j.Logger.getLogger("org.apache.hadoop.hdds.utils.db.CodecBuffer")
+ .setLevel(org.apache.log4j.Level.INFO);
+
org.apache.log4j.Logger.getLogger("org.apache.hadoop.hdds.utils.db.managed")
+ .setLevel(org.apache.log4j.Level.INFO);
+ }
+
+ @Test
+ @Timeout(value = 120, unit = TimeUnit.MINUTES)
+ public void benchmarkMixedWorkloadUnderDeletionLoad() throws Exception {
+ final String profileEvent = System.getProperty("bench.profile.event", "");
+
+ // Number of FSO buckets (in one volume) the backlog and workload are
spread across. With more than one bucket a
+ // background purge round gathers deleted dirs from several buckets, so an
ungrouped DirectoryDeletingService packs
+ // multiple buckets into one purge transaction and the apply path holds
all their write locks together; per-bucket
+ // grouping keeps each transaction single-bucket. backlogDirs below is the
TOTAL across buckets, split evenly.
+ final int numBuckets = Integer.getInteger("bench.numBuckets", 4);
+ final int backlogDirs = Integer.getInteger("bench.backlogDirs", 80);
+ final int backlogFilesPerDir =
Integer.getInteger("bench.backlogFilesPerDir", 1000);
+ final int nonEmptyEvery = Integer.getInteger("bench.backlogNonEmptyEvery",
3);
+ // Size of the stable workload dataset the read ops resolve against (never
deleted during the test).
+ final int workloadDirs = Integer.getInteger("bench.workloadDirs", 20);
+ final int workloadFilesPerDir =
Integer.getInteger("bench.workloadFilesPerDir", 100);
+ // Payload for the data-plane write/read ops and the block-bearing staged
workload files those reads hit.
+ final int fileBytes = Integer.getInteger("bench.fileBytes", 1024 * 1024);
+ final int clientThreads = Integer.getInteger("bench.clientThreads", 4);
+ final int opsPerThread = Integer.getInteger("bench.opsPerThread", 400);
+ final int pathDeletingLimit =
Integer.getInteger("bench.pathDeletingLimitPerTask", 2000);
+ final int keyDeletingLimit =
Integer.getInteger("bench.keyDeletingLimitPerTask", 40000);
+ // Phase-1 (DirectoryDeletingService) cadence. The interval is read in
MILLISECONDS (KeyManagerImpl), so 1000
+ // gives a genuine 1s gate between purge rounds: each round moves up to
pathDeletingLimit paths under one bucket
+ // write lock — the apply path this change optimizes — and gating spreads
phase 1 into a long, sampled window.
+ final int dirDeletingIntervalMs =
Integer.getInteger("bench.dirDeletingIntervalMs", 1000);
+ // Phase-2 (KeyDeletingService) is unchanged by this optimization; push
its interval past the phase-1 window so it
+ // does not run during measurement and only phase-1 contention is sampled.
+ final int blockDeletingIntervalSec =
Integer.getInteger("bench.blockDeletingIntervalSec", 600);
+ final String ratisAppenderByteLimit =
System.getProperty("bench.ratisAppenderByteLimit",
+ OMConfigKeys.OZONE_OM_RATIS_LOG_APPENDER_QUEUE_BYTE_LIMIT_DEFAULT);
+ final int perBucketBacklogDirs = Math.max(1, backlogDirs / numBuckets);
+ final int backlogFiles = numBuckets * perBucketBacklogDirs *
backlogFilesPerDir;
+
+ OzoneConfiguration conf = new OzoneConfiguration();
+ // Gate phase 1 (DirectoryDeletingService) at a genuine 1s interval so it
moves pathDeletingLimit paths per round
+ // under one apply-thread bucket write-lock hold — the path this change
optimizes — spreading phase 1 into a long,
+ // sampled window while the client workload runs. Phase 2
(KeyDeletingService) is unchanged by the optimization,
+ // so its interval is pushed past the window to keep it out of measurement.
+ conf.setInt(OMConfigKeys.OZONE_DIR_DELETING_SERVICE_INTERVAL,
dirDeletingIntervalMs);
+ conf.setTimeDuration(OZONE_BLOCK_DELETING_SERVICE_INTERVAL,
blockDeletingIntervalSec, TimeUnit.SECONDS);
+ conf.setInt(OMConfigKeys.OZONE_PATH_DELETING_LIMIT_PER_TASK,
pathDeletingLimit);
+ conf.setInt(OMConfigKeys.OZONE_KEY_DELETING_LIMIT_PER_TASK,
keyDeletingLimit);
+ conf.set(OMConfigKeys.OZONE_OM_RATIS_LOG_APPENDER_QUEUE_BYTE_LIMIT,
ratisAppenderByteLimit);
+ conf.setInt(OZONE_FS_ITERATE_BATCH_SIZE, 1000);
+
+ stripTestOnlyOverhead();
Review Comment:
Can we disable leak detection after building the mini-cluster?
`MiniOzoneClusterImpl` enables it from its static initializer, so a fresh
benchmark run currently enables it again before measurement.
##########
hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/service/TestOmMixedWorkloadUnderDeletionBench.java:
##########
@@ -0,0 +1,686 @@
+/*
+ * 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.ozone.om.service;
+
+import static
org.apache.hadoop.fs.CommonConfigurationKeysPublic.FS_DEFAULT_NAME_KEY;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_BLOCK_DELETING_SERVICE_INTERVAL;
+import static
org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_FS_ITERATE_BATCH_SIZE;
+
+import java.io.File;
+import java.io.IOException;
+import java.lang.reflect.Method;
+import java.net.URL;
+import java.net.URLClassLoader;
+import java.security.AccessController;
+import java.security.PrivilegedActionException;
+import java.security.PrivilegedExceptionAction;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Locale;
+import java.util.concurrent.Callable;
+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.AtomicBoolean;
+import org.apache.hadoop.fs.FSDataInputStream;
+import org.apache.hadoop.fs.FSDataOutputStream;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.hdds.conf.OzoneConfiguration;
+import org.apache.hadoop.hdds.protocol.StorageType;
+import org.apache.hadoop.hdds.utils.db.CodecBuffer;
+import org.apache.hadoop.hdds.utils.db.Table;
+import org.apache.hadoop.ozone.MiniOzoneCluster;
+import org.apache.hadoop.ozone.OzoneConsts;
+import org.apache.hadoop.ozone.client.BucketArgs;
+import org.apache.hadoop.ozone.client.OzoneBucket;
+import org.apache.hadoop.ozone.client.OzoneClient;
+import org.apache.hadoop.ozone.client.OzoneVolume;
+import org.apache.hadoop.ozone.om.OMConfigKeys;
+import org.apache.hadoop.ozone.om.OMMetadataManager;
+import org.apache.hadoop.ozone.om.helpers.BucketLayout;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * End-to-end benchmark reproducing the interactive-workload degradation seen
when a large background deletion backlog
+ * is being reclaimed on the same bucket that clients are actively writing to.
+ *
+ * <p>Every OM write is applied on a single serial state-machine apply thread
under a per-bucket write lock. Reclaiming
+ * a deletion backlog drives {@code PurgeDirectories} transactions through
that thread; when the deleted directories are
+ * densely populated, a single purge batch moves a large number of
sub-files/sub-dirs under one write-lock hold, so the
+ * apply thread — and the bucket write lock — is occupied for the whole batch.
Concurrent user {@code create},
+ * {@code mkdir} and {@code rename} contend for that same lock and thread, and
reads contend for the bucket read lock,
+ * so their latency degrades while the backlog drains.
+ *
+ * <p>This benchmark stages two sets under a single volume and FSO bucket: a
densely-populated <em>backlog</em> subtree
+ * that is recursively deleted and drained, and a separate stable
<em>workload</em> set that is never deleted during the
+ * run (a pre-staged dataset the read ops resolve against, plus a per-thread
scratch area the write ops create into).
+ * Sharing the bucket is intentional — deletion and the client workload
contend on the same bucket lock, while the
+ * workload always hits live paths. It then compares a mixed client workload
on that bucket — create, mkdir, rename and
+ * a data-plane file write (create + write a block), plus a data-plane file
read and the metadata read RPCs that take
+ * the bucket read lock (getFileStatus, listStatus, getBucketInfo, lookupKey)
— in two conditions:
+ * <ul>
+ * <li><b>control</b> — no deletion running, and</li>
+ * <li><b>under load</b> — the same workload while the backlog subtree is
recursively deleted and fully purged from
+ * OM in the background (both FSO phases: moved into the deletedTable,
then purged back out),</li>
+ * </ul>
+ * reporting per-operation p50/p99 latency and the under-load degradation, so
the two code versions can be compared on
+ * how much apply-thread purge work bleeds into interactive latency. It also
reports the phase-1 drain time — the
+ * {@code DirectoryDeletingService} move into the deletedTable via {@code
OMDirectoriesPurgeRequestWithFSO}, the apply
+ * path this change optimizes — separately from the full both-phase drain, so
pure apply-thread throughput can be
+ * compared alongside the interactive degradation.
+ *
+ * <p>Deletion is configured with production-representative per-task limits so
batches are large. The {@code benchmark}
+ * tag is excluded from {@code mvn test} and CI by default, so it must be
re-enabled explicitly to run on demand
+ * (rebuild the reactor first to avoid stale-class errors):
+ * <pre>
+ * mvn -pl :ozone-integration-test test -DskipShade -DskipRecon \
+ * -Dtest=TestOmMixedWorkloadUnderDeletionBench -Dgroups=benchmark
-Dexcluded-test-groups= \
+ * -Dsurefire.failIfNoSpecifiedTests=false
+ * </pre>
+ * Tunables: {@code bench.backlogDirs} (default 80), {@code
bench.backlogFilesPerDir} (default 1000),
+ * {@code bench.backlogNonEmptyEvery} (default 3 — every 3rd backlog file is
written with a block, the rest are
+ * empty so a large backlog stays cheap to stage), {@code bench.workloadDirs}
(default 20) and
+ * {@code bench.workloadFilesPerDir} (default 100) sizing the stable dataset
the read ops resolve against,
+ * {@code bench.fileBytes} (default 1 MiB) sizing the data-plane file
write/read payload (and the block-bearing
+ * staged files the reads pull), {@code bench.clientThreads} (default 4),
+ * {@code bench.opsPerThread} (default 400), {@code
bench.pathDeletingLimitPerTask} (default 2000) and
+ * {@code bench.keyDeletingLimitPerTask} (default 40000). The last two size
how much a single deletion round gathers;
+ * with the Ratis appender byte limit non-binding at these entry sizes, a
round's paths pack into one purge
+ * transaction, so raising them makes each apply move far more entries under a
single bucket write-lock hold — the
+ * regime where the apply-thread per-entry cost dominates interactive latency.
+ *
+ * <p>Adding {@code -Dbench.profile.event=<cpu|lock|wall|alloc>} profiles only
the under-load window with
+ * async-profiler, loaded reflectively from a local install whose paths must
be supplied via
+ * {@code -Dbench.profiler.jar} (the async-profiler jar) and {@code
-Dbench.profiler.lib} (its native library); the
+ * JFR is written under {@code -Dbench.profile.out}, default {@code /tmp}. For
accurate leaf frames also pass
+ * {@code -DargLine="-XX:+UnlockDiagnosticVMOptions -XX:+DebugNonSafepoints"}.
+ */
+@Tag("benchmark")
+public class TestOmMixedWorkloadUnderDeletionBench {
+
+ private static final Logger LOG =
LoggerFactory.getLogger(TestOmMixedWorkloadUnderDeletionBench.class);
+
+ private static final String OP_CREATE = "create";
+ private static final String OP_MKDIR = "mkdir";
+ private static final String OP_RENAME = "rename";
+ private static final String OP_FILEWRITE = "filewrite";
+ private static final String OP_FILEREAD = "fileread";
+ private static final String OP_GETFILESTATUS = "getfilestatus";
+ private static final String OP_LISTSTATUS = "liststatus";
+ private static final String OP_INFOBUCKET = "infobucket";
+ private static final String OP_GETKEYINFO = "getkeyinfo";
+ private static final String[] OPS =
+ {OP_CREATE, OP_MKDIR, OP_RENAME, OP_FILEWRITE, OP_FILEREAD,
+ OP_GETFILESTATUS, OP_LISTSTATUS, OP_INFOBUCKET, OP_GETKEYINFO};
+
+ // The three sandboxes share the parent /workload/bucket but are separate
subtrees; only the backlog is deleted.
+ // The deletion set (built, recursively deleted, then drained through the
apply thread).
+ private static final String BACKLOG_ROOT = "workload/bucket/backlog";
+ // The workload set — never deleted during the test: a stable pre-staged
dataset the read ops resolve against, plus
+ // a scratch area the write ops create into. Keeping this separate from the
deletion set means the concurrent
+ // operations always hit live paths while the backlog drains.
+ private static final String WORKLOAD_DATA_ROOT = "workload/bucket/data";
+ private static final String WORKLOAD_SCRATCH_ROOT =
"workload/bucket/scratch";
+
+ // Every bench.backlogNonEmptyEvery-th backlog file is written with a single
block so its KeyInfo carries a
+ // key-location list, exercising the block-metadata parse/serialize the
purge apply and flush paths hit in
+ // production. The rest are left empty because block-bearing files are far
more expensive to stage (block
+ // allocation + datanode write + commit), and a large backlog is what
actually stresses the apply thread.
+ private static final byte[] FILE_CONTENT = new byte[4];
+
+ /**
+ * Removes test-harness-only overhead that would otherwise distort the
apply/flush cost under measurement: the
+ * mini-cluster unconditionally enables {@link CodecBuffer} leak detection
(a per-allocation finalizer), and the test
+ * log config runs the {@code CodecBuffer}/managed-RocksDB loggers at
DEBUG/TRACE, which capture a full stack trace on
+ * every buffer allocation. Neither happens in a production OM running at
INFO.
+ */
+ private static void stripTestOnlyOverhead() {
+ CodecBuffer.disableLeakDetection();
+
org.apache.log4j.Logger.getLogger("org.apache.hadoop.hdds.utils.db.CodecBuffer")
+ .setLevel(org.apache.log4j.Level.INFO);
+
org.apache.log4j.Logger.getLogger("org.apache.hadoop.hdds.utils.db.managed")
+ .setLevel(org.apache.log4j.Level.INFO);
+ }
+
+ @Test
+ @Timeout(value = 120, unit = TimeUnit.MINUTES)
+ public void benchmarkMixedWorkloadUnderDeletionLoad() throws Exception {
+ final String profileEvent = System.getProperty("bench.profile.event", "");
+
+ // Number of FSO buckets (in one volume) the backlog and workload are
spread across. With more than one bucket a
+ // background purge round gathers deleted dirs from several buckets, so an
ungrouped DirectoryDeletingService packs
+ // multiple buckets into one purge transaction and the apply path holds
all their write locks together; per-bucket
+ // grouping keeps each transaction single-bucket. backlogDirs below is the
TOTAL across buckets, split evenly.
+ final int numBuckets = Integer.getInteger("bench.numBuckets", 4);
+ final int backlogDirs = Integer.getInteger("bench.backlogDirs", 80);
+ final int backlogFilesPerDir =
Integer.getInteger("bench.backlogFilesPerDir", 1000);
+ final int nonEmptyEvery = Integer.getInteger("bench.backlogNonEmptyEvery",
3);
+ // Size of the stable workload dataset the read ops resolve against (never
deleted during the test).
+ final int workloadDirs = Integer.getInteger("bench.workloadDirs", 20);
+ final int workloadFilesPerDir =
Integer.getInteger("bench.workloadFilesPerDir", 100);
+ // Payload for the data-plane write/read ops and the block-bearing staged
workload files those reads hit.
+ final int fileBytes = Integer.getInteger("bench.fileBytes", 1024 * 1024);
+ final int clientThreads = Integer.getInteger("bench.clientThreads", 4);
+ final int opsPerThread = Integer.getInteger("bench.opsPerThread", 400);
+ final int pathDeletingLimit =
Integer.getInteger("bench.pathDeletingLimitPerTask", 2000);
+ final int keyDeletingLimit =
Integer.getInteger("bench.keyDeletingLimitPerTask", 40000);
+ // Phase-1 (DirectoryDeletingService) cadence. The interval is read in
MILLISECONDS (KeyManagerImpl), so 1000
+ // gives a genuine 1s gate between purge rounds: each round moves up to
pathDeletingLimit paths under one bucket
+ // write lock — the apply path this change optimizes — and gating spreads
phase 1 into a long, sampled window.
+ final int dirDeletingIntervalMs =
Integer.getInteger("bench.dirDeletingIntervalMs", 1000);
+ // Phase-2 (KeyDeletingService) is unchanged by this optimization; push
its interval past the phase-1 window so it
+ // does not run during measurement and only phase-1 contention is sampled.
+ final int blockDeletingIntervalSec =
Integer.getInteger("bench.blockDeletingIntervalSec", 600);
+ final String ratisAppenderByteLimit =
System.getProperty("bench.ratisAppenderByteLimit",
+ OMConfigKeys.OZONE_OM_RATIS_LOG_APPENDER_QUEUE_BYTE_LIMIT_DEFAULT);
+ final int perBucketBacklogDirs = Math.max(1, backlogDirs / numBuckets);
+ final int backlogFiles = numBuckets * perBucketBacklogDirs *
backlogFilesPerDir;
+
+ OzoneConfiguration conf = new OzoneConfiguration();
+ // Gate phase 1 (DirectoryDeletingService) at a genuine 1s interval so it
moves pathDeletingLimit paths per round
+ // under one apply-thread bucket write-lock hold — the path this change
optimizes — spreading phase 1 into a long,
+ // sampled window while the client workload runs. Phase 2
(KeyDeletingService) is unchanged by the optimization,
+ // so its interval is pushed past the window to keep it out of measurement.
+ conf.setInt(OMConfigKeys.OZONE_DIR_DELETING_SERVICE_INTERVAL,
dirDeletingIntervalMs);
+ conf.setTimeDuration(OZONE_BLOCK_DELETING_SERVICE_INTERVAL,
blockDeletingIntervalSec, TimeUnit.SECONDS);
+ conf.setInt(OMConfigKeys.OZONE_PATH_DELETING_LIMIT_PER_TASK,
pathDeletingLimit);
+ conf.setInt(OMConfigKeys.OZONE_KEY_DELETING_LIMIT_PER_TASK,
keyDeletingLimit);
+ conf.set(OMConfigKeys.OZONE_OM_RATIS_LOG_APPENDER_QUEUE_BYTE_LIMIT,
ratisAppenderByteLimit);
+ conf.setInt(OZONE_FS_ITERATE_BATCH_SIZE, 1000);
+
+ stripTestOnlyOverhead();
+ MiniOzoneCluster cluster = MiniOzoneCluster.newBuilder(conf)
+ .setNumDatanodes(3)
+ .build();
+ try {
+ cluster.waitForClusterToBeReady();
+ DirectoryDeletingService dds =
cluster.getOzoneManager().getKeyManager().getDirDeletingService();
+ OMMetadataManager mm = cluster.getOzoneManager().getMetadataManager();
+
+ try (OzoneClient client = cluster.newClient()) {
+ BucketFs env = createBuckets(client, conf, numBuckets);
+ try {
+ // Stage the deletion set (recursively deleted and drained below)
and the stable workload set — a dataset the
+ // read ops resolve against plus a per-thread scratch area the write
ops create into — in every bucket, so the
+ // workload always hits live paths in the same buckets the purge is
draining. Backlog files carry only a tiny
+ // block (one key-location for the purge to parse); workload files
carry the full fileBytes payload.
+ byte[] fileContent = new byte[fileBytes];
+ for (int b = 0; b < numBuckets; b++) {
+ buildDenseTree(env.fs[b], new Path("/" + BACKLOG_ROOT),
perBucketBacklogDirs, backlogFilesPerDir,
+ nonEmptyEvery, FILE_CONTENT);
+ buildDenseTree(env.fs[b], new Path("/" + WORKLOAD_DATA_ROOT),
workloadDirs, workloadFilesPerDir,
+ nonEmptyEvery, fileContent);
+ env.fs[b].mkdirs(new Path("/" + WORKLOAD_SCRATCH_ROOT));
+ }
+ WorkloadDataset dataset = new WorkloadDataset(workloadDirs,
workloadFilesPerDir, nonEmptyEvery, fileContent);
+
+ // Control: interactive latency on the hot buckets with no deletion
running.
+ Percentiles[] control = toPercentiles(
+ startMixedWorkload(env.fs, env.volume, env.buckets, dataset,
clientThreads, opsPerThread, "control")
+ .await(), "control");
+
+ // Profile only the under-load window when
-Dbench.profile.event=<cpu|lock|wall|alloc> is set. async-profiler
+ // is loaded reflectively from -Dbench.profiler.jar /
-Dbench.profiler.lib and a JFR recording is written to
+ // -Dbench.profile.out (default /tmp). For accurate leaf frames add
+ // -DargLine="-XX:+UnlockDiagnosticVMOptions
-XX:+DebugNonSafepoints".
+ Profiler profiler = profileEvent.isEmpty() ? null : Profiler.load();
+ String profileOut = null;
+ if (profiler != null) {
+ profileOut = System.getProperty("bench.profile.out", "/tmp") +
"/prof-mixed-" + profileEvent + ".jfr";
+ profiler.start(profileEvent, profileOut);
+ }
+
+ // Run the interactive workload continuously in background threads,
sampling throughout phase 1, and stop it
+ // the moment phase 1 completes so the under-load percentiles
reflect only the optimized, contended window.
+ // Phase 1 (DirectoryDeletingService moving sub-files/sub-dirs into
the deletedTable) takes the bucket write
+ // lock on the apply thread — the path this change optimizes; every
bucket's backlog is deleted here so purge
+ // rounds span buckets. Phase 2 (KeyDeletingService draining the
deletedTable) is unchanged by the change and
+ // is kept out of the window by a long block-deleting interval, so
it is neither sampled nor waited on.
+ Table<String, ?> deletedDirTable = mm.getDeletedDirTable();
+ long movedFilesBefore = dds.getMovedFilesCount();
+ long start = System.nanoTime();
+ long drainDeadline = start + TimeUnit.SECONDS.toNanos(900);
+ RunningWorkload underLoadWl = startMixedWorkload(env.fs, env.volume,
env.buckets, dataset, clientThreads, 0,
+ "under-load");
+ for (int b = 0; b < numBuckets; b++) {
+ env.fs[b].delete(new Path("/" + BACKLOG_ROOT), true);
+ }
+ LOG.info("delete(recursive) issued on {} buckets; backlog draining
in background", numBuckets);
+
+ long phase1DrainNanos;
+ try {
+ phase1DrainNanos = awaitPhase1Drain(dds, mm, deletedDirTable,
movedFilesBefore, backlogFiles,
+ underLoadWl, start, drainDeadline);
+ } finally {
+ underLoadWl.stop();
+ if (profiler != null) {
+ profiler.stop();
+ }
+ }
+ double phase1DrainMs = phase1DrainNanos / 1_000_000.0;
+ List<List<Long>> underLoadSamples = underLoadWl.await();
+ Percentiles[] underLoad = toPercentiles(underLoadSamples,
"under-load");
+ long underLoadOps = totalOps(underLoadSamples);
+ String header = String.format(Locale.ROOT,
+ "BENCH mixed numBuckets=%d backlogFiles=%d threads=%d
opsPerThread=%d ratisByteLimit=%s underLoadOps=%d "
+ + "phase1DrainMs=%.1f",
+ numBuckets, backlogFiles, clientThreads, opsPerThread,
ratisAppenderByteLimit, underLoadOps,
+ phase1DrainMs);
+ printBenchLine(control, underLoad, header);
+ if (profileOut != null) {
+ System.out.printf(Locale.ROOT, "BENCH profile event=%s out=%s%n",
profileEvent, profileOut);
+ }
+ } finally {
+ for (FileSystem f : env.fs) {
+ org.apache.hadoop.io.IOUtils.closeStream(f);
+ }
+ }
+ }
+ } finally {
+ cluster.shutdown();
+ }
+ }
+
+ /**
+ * Blocks until phase 1 (DirectoryDeletingService) finishes: every backlog
sub-file has been moved (cheap counter)
+ * and the deletedDirTable has been fully drained of sub-dirs.
countRowsInTable scans the table, so it is only probed
+ * once the moved-files counter shows the sub-files are all moved. Phase 2
(deletedTable purge) is out of scope.
+ * Returns the drain duration in nanos, or -1 if the workload failed before
the drain completed.
+ */
+ @SuppressWarnings("checkstyle:ParameterNumber")
+ private static long awaitPhase1Drain(DirectoryDeletingService dds,
OMMetadataManager mm,
+ Table<String, ?> deletedDirTable, long movedFilesBefore, long
backlogFiles, RunningWorkload underLoadWl,
+ long start, long drainDeadline) throws Exception {
+ while (true) {
+ long moved = dds.getMovedFilesCount() - movedFilesBefore;
+ if (moved >= backlogFiles && mm.countRowsInTable(deletedDirTable) == 0) {
+ return System.nanoTime() - start;
+ }
+ if (underLoadWl.failed()) {
+ return -1;
+ }
+ if (System.nanoTime() > drainDeadline) {
+ throw new IllegalStateException("phase 1 did not drain within 900s:
moved=" + moved
+ + " expected=" + backlogFiles
+ + " deletedDirTableRows=" + mm.countRowsInTable(deletedDirTable));
+ }
+ Thread.sleep(200);
+ }
+ }
+
+ private void buildDenseTree(FileSystem fs, Path root, int dirs, int
filesPerDir, int nonEmptyEvery,
+ byte[] blockContent) throws Exception {
+ long buildStart = System.nanoTime();
+ // Staging is client/RPC-round-trip bound, not apply-thread bound, so more
concurrent creators speed it up
+ // nearly linearly; a large backlog is otherwise the long pole of the run.
Setup-only — does not affect the
+ // measured control/under-load workload.
+ ExecutorService pool =
Executors.newFixedThreadPool(Integer.getInteger("bench.stagingThreads", 48));
Review Comment:
Could we make executor cleanup unconditional? If a `Future#get()` throws,
`shutdown()` is skipped and the fixed-pool threads can keep the Maven fork
alive. The same applies to `RunningWorkload.await()`.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]