From 2f7077c06ccbd5e8e7259c6891fe98d83ec3359d Mon Sep 17 00:00:00 2001 From: Marcus Eriksson Date: Tue, 17 Feb 2015 16:20:35 +0100 Subject: [PATCH] Pick sstables to validate as late as possible with inc repairs Patch by marcuse; reviewed by yukim for CASSANDRA-8366 --- CHANGES.txt | 1 + .../cassandra/db/ColumnFamilyStore.java | 14 +++++++ .../db/compaction/CompactionManager.java | 14 ++++++- .../service/ActiveRepairService.java | 41 +++++++------------ 4 files changed, 42 insertions(+), 28 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 15a5a6125d..c3c7a19ff0 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 2.1.4 + * Pick sstables for validation as late as possible inc repairs (CASSANDRA-8366) * Fix commitlog getPendingTasks to not increment (CASSANDRA-8856) * Fix parallelism adjustment in range and secondary index queries when the first fetch does not satisfy the limit (CASSANDRA-8856) diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index 62aadf976c..e4531f22c0 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -2926,4 +2926,18 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean return new ArrayList<>(view.sstables); } }; + + public static final Function> UNREPAIRED_SSTABLES = new Function>() + { + public List apply(DataTracker.View view) + { + List sstables = new ArrayList<>(); + for (SSTableReader sstable : view.sstables) + { + if (!sstable.isRepaired()) + sstables.add(sstable); + } + return sstables; + } + }; } diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java index 68313a35a0..e54a25fae0 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java @@ -956,7 +956,19 @@ public class CompactionManager implements CompactionManagerMBean if (validator.desc.parentSessionId == null || ActiveRepairService.instance.getParentRepairSession(validator.desc.parentSessionId) == null) sstables = cfs.selectAndReference(ColumnFamilyStore.ALL_SSTABLES).refs; else - sstables = ActiveRepairService.instance.getParentRepairSession(validator.desc.parentSessionId).getAndReferenceSSTables(cfs.metadata.cfId); + { + ColumnFamilyStore.RefViewFragment refView = cfs.selectAndReference(ColumnFamilyStore.UNREPAIRED_SSTABLES); + sstables = refView.refs; + Set currentlyRepairing = ActiveRepairService.instance.currentlyRepairing(cfs.metadata.cfId, validator.desc.parentSessionId); + + if (!Sets.intersection(currentlyRepairing, Sets.newHashSet(refView.sstables)).isEmpty()) + { + logger.error("Cannot start multiple repair sessions over the same sstables"); + throw new RuntimeException("Cannot start multiple repair sessions over the same sstables"); + } + + ActiveRepairService.instance.getParentRepairSession(validator.desc.parentSessionId).addSSTables(cfs.metadata.cfId, refView.sstables); + } if (validator.gcBefore > 0) gcBefore = validator.gcBefore; diff --git a/src/java/org/apache/cassandra/service/ActiveRepairService.java b/src/java/org/apache/cassandra/service/ActiveRepairService.java index bf1cdd6360..f71cb6be66 100644 --- a/src/java/org/apache/cassandra/service/ActiveRepairService.java +++ b/src/java/org/apache/cassandra/service/ActiveRepairService.java @@ -305,38 +305,16 @@ public class ActiveRepairService public synchronized void registerParentRepairSession(UUID parentRepairSession, List columnFamilyStores, Collection> ranges) { - Map> sstablesToRepair = new HashMap<>(); - for (ColumnFamilyStore cfs : columnFamilyStores) - { - Set sstables = new HashSet<>(); - Set currentlyRepairing = currentlyRepairing(cfs.metadata.cfId); - for (SSTableReader sstable : cfs.getSSTables()) - { - if (new Bounds<>(sstable.first.getToken(), sstable.last.getToken()).intersects(ranges)) - { - if (!sstable.isRepaired()) - { - if (currentlyRepairing.contains(sstable)) - { - logger.error("Already repairing "+sstable+", can not continue."); - throw new RuntimeException("Already repairing "+sstable+", can not continue."); - } - sstables.add(sstable); - } - } - } - sstablesToRepair.put(cfs.metadata.cfId, sstables); - } - parentRepairSessions.put(parentRepairSession, new ParentRepairSession(columnFamilyStores, ranges, sstablesToRepair, System.currentTimeMillis())); + parentRepairSessions.put(parentRepairSession, new ParentRepairSession(columnFamilyStores, ranges, System.currentTimeMillis())); } - private Set currentlyRepairing(UUID cfId) + public Set currentlyRepairing(UUID cfId, UUID parentRepairSession) { Set repairing = new HashSet<>(); for (Map.Entry entry : parentRepairSessions.entrySet()) { Collection sstables = entry.getValue().sstableMap.get(cfId); - if (sstables != null) + if (sstables != null && !entry.getKey().equals(parentRepairSession)) repairing.addAll(sstables); } return repairing; @@ -419,12 +397,12 @@ public class ActiveRepairService public final Map> sstableMap; public final long repairedAt; - public ParentRepairSession(List columnFamilyStores, Collection> ranges, Map> sstables, long repairedAt) + public ParentRepairSession(List columnFamilyStores, Collection> ranges, long repairedAt) { for (ColumnFamilyStore cfs : columnFamilyStores) this.columnFamilyStores.put(cfs.metadata.cfId, cfs); this.ranges = ranges; - this.sstableMap = sstables; + this.sstableMap = new HashMap<>(); this.repairedAt = repairedAt; } @@ -452,6 +430,15 @@ public class ActiveRepairService return new Refs<>(references.build()); } + public void addSSTables(UUID cfId, Collection sstables) + { + Set existingSSTables = this.sstableMap.get(cfId); + if (existingSSTables == null) + existingSSTables = new HashSet<>(); + existingSSTables.addAll(sstables); + this.sstableMap.put(cfId, existingSSTables); + } + @Override public String toString() {