Author: catholicon
Date: Thu Nov 9 03:11:07 2017
New Revision: 1814697
URL: http://svn.apache.org/viewvc?rev=1814697&view=rev
Log:
OAK-6862: Active deletion of Lucene binaries: JMX bean, and ability to disable
automatic
* Removed scheduled execution of blob purge
* added a few logs
* add cancelability
* Added mbean to run/cancel/getStatus for purge
* Improved retry of delete from pending blob list a bit to do a resume of
cancelled purge more reliable
Added:
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBean.java
(with props)
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBeanImpl.java
(with props)
Modified:
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorFactory.java
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorTest.java
Added:
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBean.java
URL:
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBean.java?rev=1814697&view=auto
==============================================================================
---
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBean.java
(added)
+++
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBean.java
Thu Nov 9 03:11:07 2017
@@ -0,0 +1,59 @@
+/*
+ * 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.jackrabbit.oak.plugins.index.lucene;
+
+import javax.annotation.Nonnull;
+import javax.management.openmbean.CompositeData;
+
+/**
+ * MBean for starting and monitoring the progress of
+ * collection of deleted lucene index blobs.
+ *
+ * @see org.apache.jackrabbit.oak.api.jmx.RepositoryManagementMBean
+ */
+public interface ActiveDeletedBlobCollectorMBean {
+ String TYPE = "ActiveDeletedBlobCollector";
+
+ /**
+ * Initiate collection operation of deleted lucene index blobs
+ *
+ * @return the status of the operation right after it was initiated
+ */
+ @Nonnull
+ CompositeData startActiveCollection();
+
+ /**
+ * Cancel a running collection of deleted lucene index blobs operation.
+ * Does nothing if collection is not running.
+ *
+ * @return the status of the operation right after it was initiated
+ */
+ @Nonnull
+ CompositeData cancelActiveCollection();
+
+ /**
+ * Status of collection of deleted lucene index blobs.
+ *
+ * @return the status of the ongoing operation or if none the terminal
+ * status of the last operation or <em>Status not available</em> if none.
+ */
+ @Nonnull
+ CompositeData getActiveCollectionStatus();
+}
Propchange:
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBean.java
------------------------------------------------------------------------------
svn:eol-style = native
Added:
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBeanImpl.java
URL:
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBeanImpl.java?rev=1814697&view=auto
==============================================================================
---
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBeanImpl.java
(added)
+++
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBeanImpl.java
Thu Nov 9 03:11:07 2017
@@ -0,0 +1,167 @@
+/*
+ * 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.jackrabbit.oak.plugins.index.lucene;
+
+import org.apache.jackrabbit.oak.api.jmx.CheckpointMBean;
+import org.apache.jackrabbit.oak.commons.jmx.ManagementOperation;
+import
org.apache.jackrabbit.oak.plugins.index.lucene.directory.ActiveDeletedBlobCollectorFactory.ActiveDeletedBlobCollector;
+import org.apache.jackrabbit.oak.spi.blob.GarbageCollectableBlobStore;
+import org.apache.jackrabbit.oak.spi.whiteboard.Tracker;
+import org.apache.jackrabbit.oak.spi.whiteboard.Whiteboard;
+import org.apache.jackrabbit.oak.stats.Clock;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import javax.annotation.Nonnull;
+import javax.management.openmbean.CompositeData;
+import java.util.List;
+import java.util.concurrent.Callable;
+import java.util.concurrent.Executor;
+import java.util.concurrent.TimeUnit;
+
+import static com.google.common.base.Preconditions.checkNotNull;
+import static
org.apache.jackrabbit.oak.commons.jmx.ManagementOperation.Status.failed;
+import static
org.apache.jackrabbit.oak.commons.jmx.ManagementOperation.Status.initiated;
+import static org.apache.jackrabbit.oak.commons.jmx.ManagementOperation.done;
+import static
org.apache.jackrabbit.oak.commons.jmx.ManagementOperation.newManagementOperation;
+
+public class ActiveDeletedBlobCollectorMBeanImpl implements
ActiveDeletedBlobCollectorMBean {
+ private static final Logger LOG =
LoggerFactory.getLogger(ActiveDeletedBlobCollectorMBeanImpl.class);
+
+ public static final String OP_NAME = "Active lucene index blobs
collection";
+
+ /**
+ * Actively deleted blob must be deleted for at least this long (in
seconds)
+ */
+ private final long MIN_BLOB_AGE_TO_ACTIVELY_DELETE =
Long.getLong("oak.active.deletion.minAge",
+ TimeUnit.HOURS.toSeconds(24));
+
+ private final Clock clock = Clock.SIMPLE;
+
+ @Nonnull
+ private final ActiveDeletedBlobCollector activeDeletedBlobCollector;
+
+ @Nonnull
+ private Whiteboard whiteboard;
+
+ @Nonnull
+ private final GarbageCollectableBlobStore blobStore;
+
+ @Nonnull
+ private final Executor executor;
+
+
+ private ManagementOperation<Void> gcOp = done(OP_NAME, null);
+
+ /**
+ * @param activeDeletedBlobCollector deleted index blobs collector
+ * @param executor executor for running the
collection task
+ */
+ public ActiveDeletedBlobCollectorMBeanImpl(
+ @Nonnull ActiveDeletedBlobCollector activeDeletedBlobCollector,
+ @Nonnull Whiteboard whiteboard,
+ @Nonnull GarbageCollectableBlobStore blobStore,
+ @Nonnull Executor executor) {
+ this.activeDeletedBlobCollector =
checkNotNull(activeDeletedBlobCollector);
+ this.whiteboard = checkNotNull(whiteboard);
+ this.blobStore = checkNotNull(blobStore);
+ this.executor = checkNotNull(executor);
+
+ LOG.info("Active blob collector initialized with minAge: {}",
MIN_BLOB_AGE_TO_ACTIVELY_DELETE);
+ }
+
+ @Nonnull
+ @Override
+ public CompositeData startActiveCollection() {
+ if (gcOp.isDone()) {
+ long safeTimestampForDeletedBlobs =
getSafeTimestampForDeletedBlobs();
+ if (safeTimestampForDeletedBlobs == -1) {
+ return failed(OP_NAME + " couldn't be run as a safe timestamp
for" +
+ " purging lucene index blobs couldn't be
evaluated").toCompositeData();
+ }
+ gcOp = newManagementOperation(OP_NAME, () -> {
+
activeDeletedBlobCollector.purgeBlobsDeleted(safeTimestampForDeletedBlobs,
blobStore);
+ return null;
+ });
+ executor.execute(gcOp);
+ return initiated(gcOp, OP_NAME + " started").toCompositeData();
+ } else {
+ return failed(OP_NAME + " already running").toCompositeData();
+ }
+ }
+
+ @Nonnull
+ @Override
+ public CompositeData cancelActiveCollection() {
+ if (!gcOp.isDone()) {
+ executor.execute(newManagementOperation(OP_NAME, (Callable<Void>)
() -> {
+ gcOp.cancel(false);
+ activeDeletedBlobCollector.cancelBlobCollection();
+ return null;
+ }));
+ return initiated(gcOp, "Active lucene index blobs collection
cancelled").toCompositeData();
+ } else {
+ return failed(OP_NAME + " not running").toCompositeData();
+ }
+ }
+
+ @Nonnull
+ @Override
+ public CompositeData getActiveCollectionStatus() {
+ return gcOp.getStatus().toCompositeData();
+ }
+
+ private long getSafeTimestampForDeletedBlobs() {
+ long timestamp = clock.getTime() -
TimeUnit.SECONDS.toMillis(MIN_BLOB_AGE_TO_ACTIVELY_DELETE);
+
+ long minCheckpointTimestamp = getOldestCheckpointCreationTimestamp();
+
+ if (minCheckpointTimestamp == -1) {
+ return minCheckpointTimestamp;
+ }
+
+ if (minCheckpointTimestamp < timestamp) {
+ LOG.info("Oldest checkpoint timestamp ({}) is older than buffer
period ({}) for deleted blobs." +
+ " Using that instead", minCheckpointTimestamp, timestamp);
+ timestamp = minCheckpointTimestamp;
+ }
+
+ return timestamp;
+ }
+
+ private long getOldestCheckpointCreationTimestamp() {
+ Tracker<CheckpointMBean> tracker =
whiteboard.track(CheckpointMBean.class);
+
+ try {
+ List<CheckpointMBean> services = tracker.getServices();
+ if (services.size() == 1) {
+ return services.get(0).getOldestCheckpointCreationTimestamp();
+ } else if (services.isEmpty()) {
+ LOG.warn("Unable to get checkpoint mbean. No service of
required type found.");
+ return -1;
+ } else {
+ LOG.warn("Unable to get checkpoint mbean. Multiple services of
required type found.");
+ return -1;
+ }
+ } finally {
+ tracker.stop();
+ }
+ }
+}
Propchange:
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/ActiveDeletedBlobCollectorMBeanImpl.java
------------------------------------------------------------------------------
svn:eol-style = native
Modified:
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
URL:
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java?rev=1814697&r1=1814696&r2=1814697&view=diff
==============================================================================
---
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
(original)
+++
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/LuceneIndexProviderService.java
Thu Nov 9 03:11:07 2017
@@ -243,15 +243,15 @@ public class LuceneIndexProviderService
)
private static final String PROP_DISABLE_STORED_INDEX_DEFINITION =
"disableStoredIndexDefinition";
- private static final int PROP_DELETED_BLOB_COLLECTION_DEFAULT_INTERVAL =
12*60*60;
+ private static final boolean PROP_DELETED_BLOB_COLLECTION_DEFAULT_ENABLED
= true;
@Property(
- intValue = PROP_DELETED_BLOB_COLLECTION_DEFAULT_INTERVAL,
- label = "Time interval (in seconds) for actively removing deleted
index blobs from blob store",
+ boolValue = PROP_DELETED_BLOB_COLLECTION_DEFAULT_ENABLED,
+ label = "Enable actively removing deleted index blobs from blob
store",
description = "Index blobs are explicitly unique and don't require
mark-sweep type collection." +
- "This is number of seconds for scheduling clean-up. -1
would disable the functionality." +
- "Cleanup implies purging index blobs marked as deleted
earlier during some indexing cycle."
+ "This is used to enable the feature. Cleanup implies
purging index blobs marked as deleted " +
+ "earlier during some indexing cycle."
)
- private static final String
PROP_NAME_DELETED_BLOB_COLLECTION_DEFAULT_INTERVAL =
"deletedBlobsCollectionInterval";
+ private static final String
PROP_NAME_DELETED_BLOB_COLLECTION_DEFAULT_ENABLED =
"deletedBlobsCollectionEnabled";
private static final int PROP_INDEX_CLEANER_INTERVAL_DEFAULT = 10*60;
@Property(
@@ -262,13 +262,6 @@ public class LuceneIndexProviderService
)
private static final String PROP_INDEX_CLEANER_INTERVAL =
"propIndexCleanerIntervalInSecs";
- /**
- * Actively deleted blob must be deleted for at least this long (in
seconds)
- */
- final long MIN_BLOB_AGE_TO_ACTIVELY_DELETE =
Long.getLong("oak.active.deletion.minAge",
- TimeUnit.HOURS.toSeconds(24));
-
-
private static final boolean
PROP_ENABLE_SINGLE_BLOB_PER_INDEX_FILE_DEFAULT = true;
@Property(
boolValue = PROP_ENABLE_SINGLE_BLOB_PER_INDEX_FILE_DEFAULT,
@@ -749,43 +742,25 @@ public class LuceneIndexProviderService
}
private void initializeActiveBlobCollector(Whiteboard whiteboard,
Map<String, ?> config) {
- int activeDeletionInterval = PropertiesUtil.toInteger(
- config.get(PROP_NAME_DELETED_BLOB_COLLECTION_DEFAULT_INTERVAL),
- PROP_DELETED_BLOB_COLLECTION_DEFAULT_INTERVAL);
- if (activeDeletionInterval > -1 && blobStore!= null) {
+ boolean activeDeletionEnabled = PropertiesUtil.toBoolean(
+ config.get(PROP_NAME_DELETED_BLOB_COLLECTION_DEFAULT_ENABLED),
+ PROP_DELETED_BLOB_COLLECTION_DEFAULT_ENABLED);
+ if (activeDeletionEnabled && blobStore!= null) {
File blobCollectorWorkingDir = new File(indexDir, "deleted-blobs");
activeDeletedBlobCollector =
ActiveDeletedBlobCollectorFactory.newInstance(blobCollectorWorkingDir,
executorService);
- oakRegs.add(
- scheduleWithFixedDelay(whiteboard, () ->
- activeDeletedBlobCollector.purgeBlobsDeleted(
-
getSafeTimestampForDeletedBlobs(checkpointMBean),
- blobStore),
- activeDeletionInterval));
-
- log.info("Active blob collector initialized at working dir: {};
deletion interval {} seconds;" +
- "minAge: {}",
- blobCollectorWorkingDir, activeDeletionInterval,
MIN_BLOB_AGE_TO_ACTIVELY_DELETE);
+ ActiveDeletedBlobCollectorMBean bean =
+ new
ActiveDeletedBlobCollectorMBeanImpl(activeDeletedBlobCollector, whiteboard,
blobStore, executorService);
+
+ oakRegs.add(registerMBean(whiteboard,
ActiveDeletedBlobCollectorMBean.class, bean,
+ ActiveDeletedBlobCollectorMBean.TYPE, "Active lucene files
collection"));
+ log.info("Active blob collector initialized at working dir: {}",
blobCollectorWorkingDir);
} else {
activeDeletedBlobCollector =
ActiveDeletedBlobCollectorFactory.NOOP;
- log.info("Active blob collector set to NOOP. deletionInterval: {}
seconds; blobStore: {}",
- activeDeletionInterval, blobStore);
+ log.info("Active blob collector set to NOOP. enabled: {} seconds;
blobStore: {}",
+ activeDeletionEnabled, blobStore);
}
}
- private long getSafeTimestampForDeletedBlobs(CheckpointMBean
checkpointMBean) {
- long timestamp = clock.getTime() -
TimeUnit.SECONDS.toMillis(MIN_BLOB_AGE_TO_ACTIVELY_DELETE);
-
- long minCheckpointTimestamp =
checkpointMBean.getOldestCheckpointCreationTimestamp();
- if (minCheckpointTimestamp < timestamp) {
- log.info("Oldest checkpoint timestamp ({}) is older than buffer
period ({}) for deleted blobs." +
- " Using that instead", minCheckpointTimestamp, timestamp);
- timestamp = minCheckpointTimestamp;
- }
-
- return timestamp;
- }
-
-
private void registerPropertyIndexCleaner(Map<String, ?> config,
BundleContext bundleContext) {
int cleanerInterval =
PropertiesUtil.toInteger(config.get(PROP_INDEX_CLEANER_INTERVAL),
PROP_INDEX_CLEANER_INTERVAL_DEFAULT);
Modified:
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorFactory.java
URL:
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorFactory.java?rev=1814697&r1=1814696&r2=1814697&view=diff
==============================================================================
---
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorFactory.java
(original)
+++
jackrabbit/oak/trunk/oak-lucene/src/main/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorFactory.java
Thu Nov 9 03:11:07 2017
@@ -64,7 +64,10 @@ public class ActiveDeletedBlobCollectorF
* @return an instance of {@link BlobDeletionCallback} that can be
used to track deleted blobs
*/
BlobDeletionCallback getBlobDeletionCallback();
+
void purgeBlobsDeleted(long before, GarbageCollectableBlobStore
blobStore);
+
+ void cancelBlobCollection();
}
public static ActiveDeletedBlobCollector NOOP = new
ActiveDeletedBlobCollector() {
@@ -77,6 +80,11 @@ public class ActiveDeletedBlobCollectorF
public void purgeBlobsDeleted(long before, GarbageCollectableBlobStore
blobStore) {
}
+
+ @Override
+ public void cancelBlobCollection() {
+
+ }
};
public interface BlobDeletionCallback extends IndexCommitCallback {
@@ -133,6 +141,8 @@ public class ActiveDeletedBlobCollectorF
private final ExecutorService executorService;
+ private volatile boolean cancelled;
+
private static final String BLOB_FILE_PATTERN_PREFIX = "blobs-";
private static final String BLOB_FILE_PATTERN_SUFFIX = ".txt";
private static final String BLOB_FILE_PATTERN =
BLOB_FILE_PATTERN_PREFIX + "%s" + BLOB_FILE_PATTERN_SUFFIX;
@@ -166,6 +176,9 @@ public class ActiveDeletedBlobCollectorF
* @param blobStore used to purge blobs/chunks
*/
public void purgeBlobsDeleted(long before, @Nonnull
GarbageCollectableBlobStore blobStore) {
+ cancelled = false;
+ long start = clock.getTime();
+ LOG.info("Starting purge of blobs deleted before {}", before);
long numBlobsDeleted = 0;
long numChunksDeleted = 0;
@@ -185,9 +198,13 @@ public class ActiveDeletedBlobCollectorF
}
long lastCheckedBlobTimestamp = readLastCheckedBlobTimestamp();
+ long lastDeletedBlobTimestamp = lastCheckedBlobTimestamp;
String currInUseFileName = deletedBlobsFileWriter.inUseFileName;
deletedBlobsFileWriter.releaseInUseFile();
for (File deletedBlobListFile : FileUtils.listFiles(rootDirectory,
blobFileNameFilter, null)) {
+ if (cancelled) {
+ break;
+ }
if
(deletedBlobListFile.getName().equals(deletedBlobsFileWriter.inUseFileName)) {
continue;
}
@@ -202,6 +219,10 @@ public class ActiveDeletedBlobCollectorF
if (timestamp < before) {
try {
for (String deletedBlobLine :
FileUtils.readLines(deletedBlobListFile, (String) null)) {
+ if (cancelled) {
+ break;
+ }
+
String[] parsedDeletedBlobIdLine =
deletedBlobLine.split("\\|", 3);
if (parsedDeletedBlobIdLine.length != 3) {
LOG.warn("Unparseable line ({}) in file {}. It
won't be retried.",
@@ -219,6 +240,8 @@ public class ActiveDeletedBlobCollectorF
break;
}
+ lastDeletedBlobTimestamp =
Math.max(lastDeletedBlobTimestamp, blobDeletionTimestamp);
+
List<String> chunkIds =
Lists.newArrayList(blobStore.resolveChunks(deletedBlobId));
if (chunkIds.size() > 0) {
long deleted =
blobStore.countDeleteChunks(chunkIds, 0);
@@ -253,14 +276,19 @@ public class ActiveDeletedBlobCollectorF
// OAK-6314 revealed that blobs appended might not be
immediately available. So, we'd skip
// the file that was being processed when purge started -
next cycle would re-process and
// delete
- if
(!deletedBlobListFile.getName().equals(currInUseFileName) &&
!deletedBlobListFile.delete()) {
- LOG.warn("File {} couldn't be deleted while all blobs
listed in it have been purged", deletedBlobListFile);
+ if
(!deletedBlobListFile.getName().equals(currInUseFileName)) {
+ if (!deletedBlobListFile.delete()) {
+ LOG.warn("File {} couldn't be deleted while all
blobs listed in it have been purged", deletedBlobListFile);
+ } else {
+ LOG.debug("File {} deleted", deletedBlobListFile);
+ }
}
} else {
LOG.debug("Skipping {} as its timestamp is newer than {}",
deletedBlobListFile.getName(), before);
}
}
+ long startBlobTrackerSyncTime = clock.getTime();
// Synchronize deleted blob ids with the blob id tracker
try {
Closeables.close(idTempDeleteWriter, true);
@@ -274,9 +302,20 @@ public class ActiveDeletedBlobCollectorF
} catch(Exception e) {
LOG.warn("Error refreshing tracked blob ids", e);
}
+ long endBlobTrackerSyncTime = clock.getTime();
+ LOG.info("Synchronizing changes with blob tracker took {} ms",
endBlobTrackerSyncTime - startBlobTrackerSyncTime);
- LOG.info("Deleted {} blobs contained in {} chunks",
numBlobsDeleted, numChunksDeleted);
- writeOutLastCheckedBlobTimestamp(before);
+ if (cancelled) {
+ LOG.info("Deletion run cancelled by user");
+ }
+ long end = clock.getTime();
+ LOG.info("Deleted {} blobs contained in {} chunks in {} ms",
numBlobsDeleted, numChunksDeleted, end - start);
+ writeOutLastCheckedBlobTimestamp(lastDeletedBlobTimestamp);
+ }
+
+ @Override
+ public void cancelBlobCollection() {
+ cancelled = true;
}
private long readLastCheckedBlobTimestamp() {
Modified:
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorTest.java
URL:
http://svn.apache.org/viewvc/jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorTest.java?rev=1814697&r1=1814696&r2=1814697&view=diff
==============================================================================
---
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorTest.java
(original)
+++
jackrabbit/oak/trunk/oak-lucene/src/test/java/org/apache/jackrabbit/oak/plugins/index/lucene/directory/ActiveDeletedBlobCollectorTest.java
Thu Nov 9 03:11:07 2017
@@ -38,8 +38,10 @@ import java.util.Collections;
import java.util.HashSet;
import java.util.Iterator;
import java.util.List;
+import java.util.Set;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
+import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
@@ -285,6 +287,80 @@ public class ActiveDeletedBlobCollectorT
ActiveDeletedBlobCollectorFactory.NOOP, adbc);
}
+ @Test
+ public void cancellablePurge() throws Exception {
+ BlobDeletionCallback bdc = adbc.getBlobDeletionCallback();
+ for (int i = 0; i < 10; i++) {
+ String id = "Blob" + i;
+ bdc.deleted(id, Collections.singleton(id));
+ }
+ bdc.commitProgress(COMMIT_SUCCEDED);
+
+ Semaphore purgeBlocker = new Semaphore(0);
+ blobStore.callback = () -> purgeBlocker.acquireUninterruptibly();
+ Thread purgeThread = new Thread(() -> {
+ try {
+ adbc.purgeBlobsDeleted(clock.getTimeIncreasing(), blobStore);
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ }
+ });
+ purgeThread.setDaemon(true);
+ purgeBlocker.release(10);//allow 5 deletes
+ purgeThread.start();
+
+ boolean deleted5 = waitFor(5000, () ->
blobStore.deletedChunkIds.size() >= 10);
+ assertTrue("Deleted " + blobStore.deletedChunkIds.size() + " chunks",
deleted5);
+
+ adbc.cancelBlobCollection();
+ purgeBlocker.release(20);//release all that's there... this is more
than needed, btw.
+
+ boolean deleted6 = waitFor(5000, () ->
blobStore.deletedChunkIds.size() >= 12);
+ assertTrue("Haven't deleted another blob which was locked earlier.",
deleted6);
+
+ boolean cancelWorked = waitFor(5000, () -> !purgeThread.isAlive());
+ assertTrue("Cancel didn't let go of purge thread in 2 seconds",
cancelWorked);
+
+ assertTrue("Cancelling purge must return asap",
blobStore.deletedChunkIds.size() == 12);
+ }
+
+ @Test
+ public void resumeCancelledPurge() throws Exception {
+ BlobDeletionCallback bdc = adbc.getBlobDeletionCallback();
+ for (int i = 0; i < 10; i++) {
+ String id = "Blob" + i;
+ bdc.deleted(id, Collections.singleton(id));
+ }
+ bdc.commitProgress(COMMIT_SUCCEDED);
+
+ Semaphore purgeBlocker = new Semaphore(0);
+ blobStore.callback = () -> purgeBlocker.acquireUninterruptibly();
+ Thread purgeThread = new Thread(() -> {
+ try {
+ adbc.purgeBlobsDeleted(clock.getTimeIncreasing(), blobStore);
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ }
+ });
+ purgeThread.setDaemon(true);
+ purgeBlocker.release(10);//allow 5 deletes
+ purgeThread.start();
+
+ waitFor(5000, () -> blobStore.deletedChunkIds.size() >= 10);
+
+ adbc.cancelBlobCollection();
+ purgeBlocker.release(22);//release all that's there... this is more
than needed, btw.
+
+ waitFor(5000, () -> blobStore.deletedChunkIds.size() >= 12);
+
+ waitFor(5000, () -> !purgeThread.isAlive());
+
+ adbc.purgeBlobsDeleted(clock.getTimeIncreasing(), blobStore);
+
+ // Resume can re-attempt to delete already deleted blobs. Hence, the
need for for ">="
+ assertEquals("All blobs must get deleted", 20,
blobStore.deletedChunkIds.size());
+ }
+
private void verifyBlobsDeleted(String ... blobIds) throws IOException {
List<String> chunkIds = new ArrayList<>();
for (String blobId : blobIds) {
@@ -295,7 +371,8 @@ public class ActiveDeletedBlobCollectorT
}
class ChunkDeletionTrackingBlobStore implements
GarbageCollectableBlobStore {
- List<String> deletedChunkIds = Lists.newArrayList();
+ Set<String> deletedChunkIds =
com.google.common.collect.Sets.newLinkedHashSet();
+ Runnable callback = null;
volatile boolean markerChunkDeleted = false;
@Override
@@ -375,15 +452,15 @@ public class ActiveDeletedBlobCollectorT
@Override
public boolean deleteChunks(List<String> chunkIds, long
maxLastModifiedTime) throws Exception {
- deletedChunkIds.addAll(chunkIds);
setMarkerChunkDeletedFlag(chunkIds);
+ deletedChunkIds.addAll(chunkIds);
return true;
}
@Override
public long countDeleteChunks(List<String> chunkIds, long
maxLastModifiedTime) throws Exception {
- deletedChunkIds.addAll(chunkIds);
setMarkerChunkDeletedFlag(chunkIds);
+ deletedChunkIds.addAll(chunkIds);
return chunkIds.size();
}
@@ -394,6 +471,10 @@ public class ActiveDeletedBlobCollectorT
markerChunkDeleted = true;
break;
}
+
+ if (callback != null) {
+ callback.run();
+ }
}
}
}
@@ -403,4 +484,23 @@ public class ActiveDeletedBlobCollectorT
return Iterators.forArray(blobId + "-1", blobId + "-2");
}
}
+
+ private interface Condition {
+ boolean evaluate();
+ }
+
+ private boolean waitFor(long timeout, Condition c)
+ throws InterruptedException {
+ long end = System.currentTimeMillis() + timeout;
+ long remaining = end - System.currentTimeMillis();
+ while (remaining > 0) {
+ if (c.evaluate()) {
+ return true;
+ }
+
+ Thread.sleep(100);//The constant is exaggerated
+ remaining = end - System.currentTimeMillis();
+ }
+ return c.evaluate();
+ }
}