Repository: cassandra Updated Branches: refs/heads/trunk c7b77eb60 -> 173dbab24
Add more repair metrics Patch by marcuse; reviewed by Chris Lohfink for CASSANDRA-13531 Project: http://git-wip-us.apache.org/repos/asf/cassandra/repo Commit: http://git-wip-us.apache.org/repos/asf/cassandra/commit/173dbab2 Tree: http://git-wip-us.apache.org/repos/asf/cassandra/tree/173dbab2 Diff: http://git-wip-us.apache.org/repos/asf/cassandra/diff/173dbab2 Branch: refs/heads/trunk Commit: 173dbab24b9acba46279ca3477f080e57f52636c Parents: c7b77eb Author: Marcus Eriksson <[email protected]> Authored: Mon May 15 09:09:12 2017 +0200 Committer: Marcus Eriksson <[email protected]> Committed: Mon May 15 16:42:36 2017 +0200 ---------------------------------------------------------------------- CHANGES.txt | 1 + doc/source/operating/metrics.rst | 10 ++++ .../db/compaction/CompactionManager.java | 26 +++++++-- .../cassandra/metrics/KeyspaceMetrics.java | 22 +++++++- .../apache/cassandra/metrics/TableMetrics.java | 57 ++++++++++++++++++++ .../apache/cassandra/repair/LocalSyncTask.java | 2 + .../apache/cassandra/repair/RemoteSyncTask.java | 1 + .../apache/cassandra/repair/RepairRunnable.java | 4 +- .../org/apache/cassandra/repair/SyncTask.java | 10 ++++ 9 files changed, 126 insertions(+), 7 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/cassandra/blob/173dbab2/CHANGES.txt ---------------------------------------------------------------------- diff --git a/CHANGES.txt b/CHANGES.txt index 214f567..1191086 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.0 + * Add more repair metrics (CASSANDRA-13531) * RangeStreamer should be smarter when picking endpoints for streaming (CASSANDRA-4650) * Avoid rewrapping an exception thrown for cache load functions (CASSANDRA-13367) * Log time elapsed for each incremental repair phase (CASSANDRA-13498) http://git-wip-us.apache.org/repos/asf/cassandra/blob/173dbab2/doc/source/operating/metrics.rst ---------------------------------------------------------------------- diff --git a/doc/source/operating/metrics.rst b/doc/source/operating/metrics.rst index dcf82bc..bfa0527 100644 --- a/doc/source/operating/metrics.rst +++ b/doc/source/operating/metrics.rst @@ -132,6 +132,14 @@ SpeculativeInsufficientReplicas Counter Number of speculative ret SpeculativeSampleLatencyNanos Gauge<Long> Number of nanoseconds to wait before speculation is attempted. Value may be statically configured or updated periodically based on coordinator latency. WaitingOnFreeMemtableSpace Histogram Histogram of time spent waiting for free memtable space, either on- or off-heap. DroppedMutations Counter Number of dropped mutations on this table. +AnticompactionTime Timer Time spent anticompacting before a consistent repair. +ValidationTime Timer Time spent doing validation compaction during repair. +SyncTime Timer Time spent doing streaming during repair. +BytesValidated Histogram Histogram over the amount of bytes read during validation. +PartitionsValidated Histogram Histogram over the number of partitions read during validation. +BytesAnticompacted Counter How many bytes we anticompacted. +BytesMutatedAnticompaction Counter How many bytes we avoided anticompacting because the sstable was fully contained in the repaired range. +MutatedAnticompactionGauge Gauge<Double> Ratio of bytes mutated vs total bytes repaired. ======================================= ============== =========== Keyspace Metrics @@ -154,6 +162,8 @@ Name Type Description ======================================= ============== =========== WriteFailedIdeaCL Counter Number of writes that failed to achieve the configured ideal consistency level or 0 if none is configured IdealCLWriteLatency Latency Coordinator latency of writes at the configured ideal consistency level. No values are recorded if ideal consistency level is not configured +RepairTime Timer Total time spent as repair coordinator. +RepairPrepareTime Timer Total time spent preparing for repair. ======================================= ============== =========== ThreadPool Metrics http://git-wip-us.apache.org/repos/asf/cassandra/blob/173dbab2/src/java/org/apache/cassandra/db/compaction/CompactionManager.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java index 6f311a2..8efeeb0 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java @@ -66,6 +66,7 @@ import org.apache.cassandra.io.sstable.metadata.MetadataCollector; import org.apache.cassandra.io.sstable.metadata.StatsMetadata; import org.apache.cassandra.io.util.FileUtils; import org.apache.cassandra.metrics.CompactionMetrics; +import org.apache.cassandra.metrics.TableMetrics; import org.apache.cassandra.repair.Validator; import org.apache.cassandra.schema.CompactionParams.TombstoneOption; import org.apache.cassandra.service.ActiveRepairService; @@ -576,7 +577,10 @@ public class CompactionManager implements CompactionManagerMBean { protected void runMayThrow() throws Exception { - performAnticompaction(cfs, ranges, sstables, txn, ActiveRepairService.UNREPAIRED_SSTABLE, sessionId, sessionId); + try (TableMetrics.TableTimer.Context ctx = cfs.metric.anticompactionTime.time()) + { + performAnticompaction(cfs, ranges, sstables, txn, ActiveRepairService.UNREPAIRED_SSTABLE, sessionId, sessionId); + } } }; @@ -669,6 +673,7 @@ public class CompactionManager implements CompactionManagerMBean sstableIterator.remove(); } } + cfs.metric.bytesMutatedAnticompaction.inc(SSTableReader.getTotalBytes(mutatedRepairStatuses)); cfs.getTracker().notifySSTableRepairedStatusChanged(mutatedRepairStatusToNotify); txn.cancel(Sets.union(nonAnticompacting, mutatedRepairStatuses)); validatedForRepair.release(Sets.union(nonAnticompacting, mutatedRepairStatuses)); @@ -920,7 +925,7 @@ public class CompactionManager implements CompactionManagerMBean { public Object call() throws IOException { - try + try (TableMetrics.TableTimer.Context c = cfStore.metric.validationTime.time()) { doValidationCompaction(cfStore, validator); } @@ -1345,6 +1350,7 @@ public class CompactionManager implements CompactionManagerMBean // We blindly assume that a partition is evenly distributed on all sstables for now. MerkleTrees tree = createMerkleTrees(sstables, validator.desc.ranges, cfs); long start = System.nanoTime(); + long partitionCount = 0; try (AbstractCompactionStrategy.ScannerList scanners = cfs.getCompactionStrategyManager().getScanners(sstables, validator.desc.ranges); ValidationCompactionController controller = new ValidationCompactionController(cfs, gcBefore); CompactionIterator ci = new ValidationCompactionIterator(scanners.scanners, controller, nowInSec, metrics)) @@ -1358,6 +1364,7 @@ public class CompactionManager implements CompactionManagerMBean try (UnfilteredRowIterator partition = ci.next()) { validator.add(partition); + partitionCount++; } } validator.complete(); @@ -1370,12 +1377,21 @@ public class CompactionManager implements CompactionManagerMBean // is done). cfs.clearSnapshot(snapshotName); } + cfs.metric.partitionsValidated.update(partitionCount); } - + long estimatedTotalBytes = 0; + for (SSTableReader sstable : sstables) + { + for (Pair<Long, Long> positionsForRanges : sstable.getPositionsForRanges(validator.desc.ranges)) + estimatedTotalBytes += positionsForRanges.right - positionsForRanges.left; + } + cfs.metric.bytesValidated.update(estimatedTotalBytes); if (logger.isDebugEnabled()) { long duration = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - start); - logger.debug("Validation finished in {} msec, for {}", + logger.debug("Validation of {} partitions (~{}) finished in {} msec, for {}", + partitionCount, + FBUtilities.prettyPrintMemory(estimatedTotalBytes), duration, validator.desc); } @@ -1490,7 +1506,7 @@ public class CompactionManager implements CompactionManagerMBean // repairedAt values for these, we still avoid anti-compacting already repaired sstables, as we currently don't // make use of any actual repairedAt value and splitting up sstables just for that is not worth it at this point. Set<SSTableReader> unrepairedSSTables = sstables.stream().filter((s) -> !s.isRepaired()).collect(Collectors.toSet()); - + cfs.metric.bytesAnticompacted.inc(SSTableReader.getTotalBytes(unrepairedSSTables)); Collection<Collection<SSTableReader>> groupedSSTables = cfs.getCompactionStrategyManager().groupSSTablesForAntiCompaction(unrepairedSSTables); // iterate over sstables to check if the repaired / unrepaired ranges intersect them. http://git-wip-us.apache.org/repos/asf/cassandra/blob/173dbab2/src/java/org/apache/cassandra/metrics/KeyspaceMetrics.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/metrics/KeyspaceMetrics.java b/src/java/org/apache/cassandra/metrics/KeyspaceMetrics.java index 3c6b604..affb372 100644 --- a/src/java/org/apache/cassandra/metrics/KeyspaceMetrics.java +++ b/src/java/org/apache/cassandra/metrics/KeyspaceMetrics.java @@ -22,7 +22,6 @@ import java.util.Set; import com.codahale.metrics.Counter; import com.codahale.metrics.Gauge; import com.codahale.metrics.Histogram; -import com.codahale.metrics.Meter; import com.codahale.metrics.Timer; import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.Keyspace; @@ -103,6 +102,20 @@ public class KeyspaceMetrics public final Counter speculativeFailedRetries; /** Needed to speculate, but didn't have enough replicas **/ public final Counter speculativeInsufficientReplicas; + /** total time spent as a repair coordinator */ + public final Timer repairTime; + /** total time spent preparing for repair */ + public final Timer repairPrepareTime; + /** Time spent anticompacting */ + public final Timer anticompactionTime; + /** total time spent creating merkle trees */ + public final Timer validationTime; + /** total time spent syncing data after repair */ + public final Timer repairSyncTime; + /** histogram over the number of bytes we have validated */ + public final Histogram bytesValidated; + /** histogram over the number of partitions we have validated */ + public final Histogram partitionsValidated; public final MetricNameFactory factory; private Keyspace keyspace; @@ -272,6 +285,13 @@ public class KeyspaceMetrics return metric.speculativeInsufficientReplicas.getCount(); } }); + repairTime = Metrics.timer(factory.createMetricName("RepairTime")); + repairPrepareTime = Metrics.timer(factory.createMetricName("RepairPrepareTime")); + anticompactionTime = Metrics.timer(factory.createMetricName("AntiCompactionTime")); + validationTime = Metrics.timer(factory.createMetricName("ValidationTime")); + repairSyncTime = Metrics.timer(factory.createMetricName("RepairSyncTime")); + partitionsValidated = Metrics.histogram(factory.createMetricName("PartitionsValidated"), false); + bytesValidated = Metrics.histogram(factory.createMetricName("BytesValidated"), false); } /** http://git-wip-us.apache.org/repos/asf/cassandra/blob/173dbab2/src/java/org/apache/cassandra/metrics/TableMetrics.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/metrics/TableMetrics.java b/src/java/org/apache/cassandra/metrics/TableMetrics.java index c4b0000..9acdaec 100644 --- a/src/java/org/apache/cassandra/metrics/TableMetrics.java +++ b/src/java/org/apache/cassandra/metrics/TableMetrics.java @@ -146,6 +146,22 @@ public class TableMetrics public final LatencyMetrics casCommit; /** percent of the data that is repaired */ public final Gauge<Double> percentRepaired; + /** time spent anticompacting data before participating in a consistent repair */ + public final TableTimer anticompactionTime; + /** time spent creating merkle trees */ + public final TableTimer validationTime; + /** time spent syncing data in a repair */ + public final TableTimer syncTime; + /** approximate number of bytes read while creating merkle trees */ + public final TableHistogram bytesValidated; + /** number of partitions read creating merkle trees */ + public final TableHistogram partitionsValidated; + /** number of bytes read while doing anticompaction */ + public final Counter bytesAnticompacted; + /** number of bytes where the whole sstable was contained in a repairing range so that we only mutated the repair status */ + public final Counter bytesMutatedAnticompaction; + /** ratio of how much we anticompact vs how much we could mutate the repair status*/ + public final Gauge<Double> mutatedAnticompactionGauge; public final Timer coordinatorReadLatency; public final Timer coordinatorScanLatency; @@ -706,6 +722,23 @@ public class TableMetrics casPrepare = new LatencyMetrics(factory, "CasPrepare", cfs.keyspace.metric.casPrepare); casPropose = new LatencyMetrics(factory, "CasPropose", cfs.keyspace.metric.casPropose); casCommit = new LatencyMetrics(factory, "CasCommit", cfs.keyspace.metric.casCommit); + + anticompactionTime = createTableTimer("AnticompactionTime", cfs.keyspace.metric.anticompactionTime); + validationTime = createTableTimer("ValidationTime", cfs.keyspace.metric.validationTime); + syncTime = createTableTimer("SyncTime", cfs.keyspace.metric.repairSyncTime); + + bytesValidated = createTableHistogram("BytesValidated", cfs.keyspace.metric.bytesValidated, false); + partitionsValidated = createTableHistogram("PartitionsValidated", cfs.keyspace.metric.partitionsValidated, false); + bytesAnticompacted = createTableCounter("BytesAnticompacted"); + bytesMutatedAnticompaction = createTableCounter("BytesMutatedAnticompaction"); + mutatedAnticompactionGauge = createTableGauge("MutatedAnticompactionGauge", () -> + { + double bytesMutated = bytesMutatedAnticompaction.getCount(); + double bytesAnticomp = bytesAnticompacted.getCount(); + if (bytesAnticomp + bytesMutated > 0) + return bytesMutated / (bytesAnticomp + bytesMutated); + return 0.0; + }); } public void updateSSTableIterated(int count) @@ -918,6 +951,30 @@ public class TableMetrics timer.update(i, unit); } } + + public Context time() + { + return new Context(all); + } + + public static class Context implements AutoCloseable + { + private final long start; + private final Timer [] all; + + private Context(Timer [] all) + { + this.all = all; + start = System.nanoTime(); + } + + public void close() + { + long duration = System.nanoTime() - start; + for (Timer t : all) + t.update(duration, TimeUnit.NANOSECONDS); + } + } } static class TableMetricNameFactory implements MetricNameFactory http://git-wip-us.apache.org/repos/asf/cassandra/blob/173dbab2/src/java/org/apache/cassandra/repair/LocalSyncTask.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/repair/LocalSyncTask.java b/src/java/org/apache/cassandra/repair/LocalSyncTask.java index 5bb66dc..343950b 100644 --- a/src/java/org/apache/cassandra/repair/LocalSyncTask.java +++ b/src/java/org/apache/cassandra/repair/LocalSyncTask.java @@ -129,10 +129,12 @@ public class LocalSyncTask extends SyncTask implements StreamEventHandler logger.info("{} {}", previewKind.logPrefix(desc.sessionId), message); Tracing.traceRepair(message); set(stat.withSummaries(result.createSummaries())); + finished(); } public void onFailure(Throwable t) { setException(t); + finished(); } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/173dbab2/src/java/org/apache/cassandra/repair/RemoteSyncTask.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/repair/RemoteSyncTask.java b/src/java/org/apache/cassandra/repair/RemoteSyncTask.java index a26a392..6cc786e 100644 --- a/src/java/org/apache/cassandra/repair/RemoteSyncTask.java +++ b/src/java/org/apache/cassandra/repair/RemoteSyncTask.java @@ -69,5 +69,6 @@ public class RemoteSyncTask extends SyncTask { setException(new RepairException(desc, previewKind, String.format("Sync failed between %s and %s", r1.endpoint, r2.endpoint))); } + finished(); } } http://git-wip-us.apache.org/repos/asf/cassandra/blob/173dbab2/src/java/org/apache/cassandra/repair/RepairRunnable.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/repair/RepairRunnable.java b/src/java/org/apache/cassandra/repair/RepairRunnable.java index 354d8b5..f327757 100644 --- a/src/java/org/apache/cassandra/repair/RepairRunnable.java +++ b/src/java/org/apache/cassandra/repair/RepairRunnable.java @@ -34,10 +34,12 @@ import org.apache.commons.lang3.time.DurationFormatUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import com.codahale.metrics.Timer; import org.apache.cassandra.concurrent.JMXConfigurableThreadPoolExecutor; import org.apache.cassandra.concurrent.NamedThreadFactory; import org.apache.cassandra.repair.consistent.SyncStatSummary; import org.apache.cassandra.schema.SchemaConstants; +import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.cql3.QueryOptions; import org.apache.cassandra.cql3.QueryProcessor; import org.apache.cassandra.cql3.UntypedResultSet; @@ -226,7 +228,7 @@ public class RepairRunnable extends WrappedRunnable implements ProgressEventNoti } long repairedAt; - try + try (Timer.Context ctx = Keyspace.open(keyspace).metric.repairPrepareTime.time()) { ActiveRepairService.instance.prepareForRepair(parentSession, FBUtilities.getBroadcastAddress(), allNeighbors, options, columnFamilyStores); repairedAt = ActiveRepairService.instance.getParentRepairSession(parentSession).getRepairedAt(); http://git-wip-us.apache.org/repos/asf/cassandra/blob/173dbab2/src/java/org/apache/cassandra/repair/SyncTask.java ---------------------------------------------------------------------- diff --git a/src/java/org/apache/cassandra/repair/SyncTask.java b/src/java/org/apache/cassandra/repair/SyncTask.java index 6f3157a..f7cf5f1 100644 --- a/src/java/org/apache/cassandra/repair/SyncTask.java +++ b/src/java/org/apache/cassandra/repair/SyncTask.java @@ -18,11 +18,13 @@ package org.apache.cassandra.repair; import java.util.List; +import java.util.concurrent.TimeUnit; import com.google.common.util.concurrent.AbstractFuture; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.db.Keyspace; import org.apache.cassandra.dht.Range; import org.apache.cassandra.dht.Token; import org.apache.cassandra.streaming.PreviewKind; @@ -43,6 +45,7 @@ public abstract class SyncTask extends AbstractFuture<SyncStat> implements Runna protected final PreviewKind previewKind; protected volatile SyncStat stat; + protected long startTime = Long.MIN_VALUE; public SyncTask(RepairJobDesc desc, TreeResponse r1, TreeResponse r2, PreviewKind previewKind) { @@ -57,6 +60,7 @@ public abstract class SyncTask extends AbstractFuture<SyncStat> implements Runna */ public void run() { + startTime = System.currentTimeMillis(); // compare trees, and collect differences List<Range<Token>> differences = MerkleTrees.difference(r1.trees, r2.trees); @@ -83,5 +87,11 @@ public abstract class SyncTask extends AbstractFuture<SyncStat> implements Runna return stat; } + protected void finished() + { + if (startTime != Long.MIN_VALUE) + Keyspace.open(desc.keyspace).getColumnFamilyStore(desc.columnFamily).metric.syncTime.update(System.currentTimeMillis() - startTime, TimeUnit.MILLISECONDS); + } + protected abstract void startSync(List<Range<Token>> differences); } --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
