diff --git a/CHANGES.txt b/CHANGES.txt index dbfa272012..8be989c219 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.0-beta5 + * Prevent parent repair sessions leak (CASSANDRA-16446) * Fix timestamp issue in SinglePartitionSliceCommandTest testPartitionD…eletionRowDeletionTie (CASSANDRA-16443) * Promote protocol V5 out of beta (CASSANDRA-14973) * Fix incorrect encoding for strings can be UTF8 (CASSANDRA-16429) diff --git a/src/java/org/apache/cassandra/repair/RepairRunnable.java b/src/java/org/apache/cassandra/repair/RepairRunnable.java index 5d8e945ce2..793d2f2163 100644 --- a/src/java/org/apache/cassandra/repair/RepairRunnable.java +++ b/src/java/org/apache/cassandra/repair/RepairRunnable.java @@ -393,15 +393,29 @@ public class RepairRunnable implements Runnable, ProgressEventNotifier { if (options.isPreview()) { - previewRepair(parentSession, creationTimeMillis, neighborsAndRanges.filterCommonRanges(keyspace, cfnames), cfnames); + previewRepair(parentSession, + creationTimeMillis, + neighborsAndRanges.filterCommonRanges(keyspace, cfnames), + neighborsAndRanges.participants, + cfnames); } else if (options.isIncremental()) { - incrementalRepair(parentSession, creationTimeMillis, traceState, neighborsAndRanges, cfnames); + incrementalRepair(parentSession, + creationTimeMillis, + traceState, + neighborsAndRanges, + neighborsAndRanges.participants, + cfnames); } else { - normalRepair(parentSession, creationTimeMillis, traceState, neighborsAndRanges.filterCommonRanges(keyspace, cfnames), cfnames); + normalRepair(parentSession, + creationTimeMillis, + traceState, + neighborsAndRanges.filterCommonRanges(keyspace, cfnames), + neighborsAndRanges.participants, + cfnames); } } @@ -409,6 +423,7 @@ public class RepairRunnable implements Runnable, ProgressEventNotifier long startTime, TraceState traceState, List commonRanges, + Set preparedEndpoints, String... cfnames) { @@ -447,13 +462,22 @@ public class RepairRunnable implements Runnable, ProgressEventNotifier return Futures.immediateFuture(null); } }, MoreExecutors.directExecutor()); - Futures.addCallback(repairResult, new RepairCompleteCallback(parentSession, successfulRanges, startTime, traceState, hasFailure, executor), MoreExecutors.directExecutor()); + Futures.addCallback(repairResult, + new RepairCompleteCallback(parentSession, + successfulRanges, + preparedEndpoints, + startTime, + traceState, + hasFailure, + executor), + MoreExecutors.directExecutor()); } private void incrementalRepair(UUID parentSession, long startTime, TraceState traceState, NeighborsAndRanges neighborsAndRanges, + Set preparedEndpoints, String... cfnames) { // the local node also needs to be included in the set of participants, since coordinator sessions aren't persisted @@ -474,12 +498,15 @@ public class RepairRunnable implements Runnable, ProgressEventNotifier { ranges.addAll(range); } - Futures.addCallback(repairResult, new RepairCompleteCallback(parentSession, ranges, startTime, traceState, hasFailure, executor), MoreExecutors.directExecutor()); + Futures.addCallback(repairResult, + new RepairCompleteCallback(parentSession, ranges, preparedEndpoints, startTime, traceState, hasFailure, executor), + MoreExecutors.directExecutor()); } private void previewRepair(UUID parentSession, long startTime, List commonRanges, + Set preparedEndpoints, String... cfnames) { @@ -521,6 +548,7 @@ public class RepairRunnable implements Runnable, ProgressEventNotifier notification(message); success("Repair preview completed successfully"); + ActiveRepairService.instance.cleanUp(parentSession, preparedEndpoints); } catch (Throwable t) { @@ -669,6 +697,7 @@ public class RepairRunnable implements Runnable, ProgressEventNotifier { final UUID parentSession; final Collection> successfulRanges; + final Set preparedEndpoints; final long startTime; final TraceState traceState; final AtomicBoolean hasFailure; @@ -676,6 +705,7 @@ public class RepairRunnable implements Runnable, ProgressEventNotifier public RepairCompleteCallback(UUID parentSession, Collection> successfulRanges, + Set preparedEndpoints, long startTime, TraceState traceState, AtomicBoolean hasFailure, @@ -683,6 +713,7 @@ public class RepairRunnable implements Runnable, ProgressEventNotifier { this.parentSession = parentSession; this.successfulRanges = successfulRanges; + this.preparedEndpoints = preparedEndpoints; this.startTime = startTime; this.traceState = traceState; this.hasFailure = hasFailure; @@ -699,6 +730,7 @@ public class RepairRunnable implements Runnable, ProgressEventNotifier else { success("Repair completed successfully"); + ActiveRepairService.instance.cleanUp(parentSession, preparedEndpoints); } executor.shutdownNow(); } diff --git a/src/java/org/apache/cassandra/service/ActiveRepairService.java b/src/java/org/apache/cassandra/service/ActiveRepairService.java index 2cdc794473..58587beb1e 100644 --- a/src/java/org/apache/cassandra/service/ActiveRepairService.java +++ b/src/java/org/apache/cassandra/service/ActiveRepairService.java @@ -64,6 +64,7 @@ import org.apache.cassandra.locator.InetAddressAndPort; import org.apache.cassandra.locator.TokenMetadata; import org.apache.cassandra.metrics.RepairMetrics; import org.apache.cassandra.net.RequestCallback; +import org.apache.cassandra.net.Verb; import org.apache.cassandra.net.Message; import org.apache.cassandra.net.MessagingService; import org.apache.cassandra.repair.CommonRange; @@ -77,6 +78,7 @@ import org.apache.cassandra.repair.consistent.admin.PendingStats; import org.apache.cassandra.repair.consistent.admin.RepairStats; import org.apache.cassandra.repair.consistent.RepairedState; import org.apache.cassandra.repair.consistent.admin.SchemaArgsParser; +import org.apache.cassandra.repair.messages.CleanupMessage; import org.apache.cassandra.repair.messages.PrepareMessage; import org.apache.cassandra.repair.messages.RepairMessage; import org.apache.cassandra.repair.messages.RepairOption; @@ -304,6 +306,12 @@ public class ActiveRepairService implements IEndpointStateChangeSubscriber, IFai return stats; } + @Override + public int parentRepairSessionsCount() + { + return parentRepairSessions.size(); + } + /** * Requests repairs for the given keyspace and column families. * @@ -597,6 +605,49 @@ public class ActiveRepairService implements IEndpointStateChangeSubscriber, IFai return parentRepairSession; } + /** + * Send Verb.CLEANUP_MSG to the given endpoints. This results in removing parent session object from the + * endpoint's cache. + * This method does not throw an exception in case of a messaging failure. + */ + public void cleanUp(UUID parentRepairSession, Set endpoints) + { + for (InetAddressAndPort endpoint : endpoints) + { + try + { + if (FailureDetector.instance.isAlive(endpoint)) + { + CleanupMessage message = new CleanupMessage(parentRepairSession); + Message msg = Message.out(Verb.CLEANUP_MSG, message); + + RequestCallback loggingCallback = new RequestCallback() + { + @Override + public void onResponse(Message msg) + { + logger.trace("Successfully cleaned up {} parent repair session on {}.", parentRepairSession, endpoint); + } + + @Override + public void onFailure(InetAddressAndPort from, RequestFailureReason failureReason) + { + logger.debug("Failed to clean up parent repair session {} on {}. The uncleaned sessions will " + + "be removed on a node restart. This should not be a problem unless you see thousands " + + "of messages like this.", parentRepairSession, endpoint); + } + }; + + MessagingService.instance().sendWithCallback(msg, endpoint, loggingCallback); + } + } + catch (Exception exc) + { + logger.warn("Failed to send a clean up message to {}", endpoint, exc); + } + } + } + private void failRepair(UUID parentRepairSession, String errorMsg) { removeParentRepairSession(parentRepairSession); throw new RuntimeException(errorMsg); diff --git a/src/java/org/apache/cassandra/service/ActiveRepairServiceMBean.java b/src/java/org/apache/cassandra/service/ActiveRepairServiceMBean.java index 8cffecc16d..b68cb6f507 100644 --- a/src/java/org/apache/cassandra/service/ActiveRepairServiceMBean.java +++ b/src/java/org/apache/cassandra/service/ActiveRepairServiceMBean.java @@ -41,4 +41,13 @@ public interface ActiveRepairServiceMBean public List getRepairStats(List schemaArgs, String rangeString); public List getPendingStats(List schemaArgs, String rangeString); public List cleanupPending(List schemaArgs, String rangeString, boolean force); + + /** + * Each ongoing repair (incremental and non-incremental) is represented by a + * {@link ActiveRepairService.ParentRepairSession} entry in the {@link ActiveRepairService} cache. + * Returns the current number of ongoing repairs (the current number of cached entries). + * + * @return current size of the internal cache holding {@link ActiveRepairService.ParentRepairSession} instances + */ + int parentRepairSessionsCount(); }