diff --git a/build.sbt b/build.sbt index 6c1b21b4..30950a29 100644 --- a/build.sbt +++ b/build.sbt @@ -50,6 +50,15 @@ libraryDependencies ++= { ) } +libraryDependencies ++= { + CrossVersion.partialVersion(scalaVersion.value) match { + case Some((2, major)) if major <= 12 => + Seq() + case _ => + Seq("org.scala-lang.modules" %% "scala-parallel-collections" % "1.0.4") + } +} + // Set the artifact names. artifactName := { (scalaVersion: ScalaVersion, module: ModuleID, artifact: Artifact) => artifact.`type` match { diff --git a/src/main/scala-2.11/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala b/src/main/scala-2.11/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala new file mode 100644 index 00000000..2af3dc89 --- /dev/null +++ b/src/main/scala-2.11/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala @@ -0,0 +1,7 @@ +package com.comcast.xfinity.sirius.uberstore.segmented + +import scala.collection.parallel.ParSeq + +object ParallelHelpers { + def parallelize[T](seq: Seq[T]): ParSeq[T] = seq.par +} diff --git a/src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala b/src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala new file mode 100644 index 00000000..2af3dc89 --- /dev/null +++ b/src/main/scala-2.12/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala @@ -0,0 +1,7 @@ +package com.comcast.xfinity.sirius.uberstore.segmented + +import scala.collection.parallel.ParSeq + +object ParallelHelpers { + def parallelize[T](seq: Seq[T]): ParSeq[T] = seq.par +} diff --git a/src/main/scala-2.13/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala b/src/main/scala-2.13/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala new file mode 100644 index 00000000..523f15bb --- /dev/null +++ b/src/main/scala-2.13/com/comcast/xfinity/sirius/uberstore/segmented/ParallelHelpers.scala @@ -0,0 +1,8 @@ +package com.comcast.xfinity.sirius.uberstore.segmented + +import scala.collection.parallel.CollectionConverters._ +import scala.collection.parallel.ParSeq + +object ParallelHelpers { + def parallelize[T](seq: Seq[T]): ParSeq[T] = seq.par +} diff --git a/src/main/scala/com/comcast/xfinity/sirius/api/SiriusConfiguration.scala b/src/main/scala/com/comcast/xfinity/sirius/api/SiriusConfiguration.scala index 38af1b78..e3f5b857 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/api/SiriusConfiguration.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/api/SiriusConfiguration.scala @@ -308,6 +308,16 @@ object SiriusConfiguration { * Maximum akka message size in KB. Default is 1024. Type is Integer. */ final val MAX_AKKA_MESSAGE_SIZE_KB = "sirius.akka.maximum-frame-size-kb" + + /** + * Whether or not to bootstrap the log in parallel, only applies to segmented uberstores + */ + final val LOG_PARALLEL_ENABLED = "sirius.log.parallel-enabled" + + /** + * Whether to skip checksum validation when reading events from the log + */ + final val LOG_SKIP_CHECKSUM_VALIDATION = "sirius.log.skip-checksum-validation" } /** diff --git a/src/main/scala/com/comcast/xfinity/sirius/api/impl/state/StateSup.scala b/src/main/scala/com/comcast/xfinity/sirius/api/impl/state/StateSup.scala index cc9ad3c1..c9635279 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/api/impl/state/StateSup.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/api/impl/state/StateSup.scala @@ -119,22 +119,15 @@ class StateSup(requestHandler: RequestHandler, bootstrapTime = Some(0L) case _ => + val parallel = config.getProp(SiriusConfiguration.LOG_PARALLEL_ENABLED, default = false) val start = System.currentTimeMillis logger.info("Beginning SiriusLog replay at {}", start) requestHandler.onBootstrapStarting() - siriusLog.foreach( - orderedEvent => - try { - orderedEvent.request match { - case Put(key, body) => requestHandler.handlePut(orderedEvent.sequence, key, body) - case Delete(key) => requestHandler.handleDelete(orderedEvent.sequence, key) - } - } catch { - case rte: RuntimeException => - eventReplayFailureCount += 1 - logger.error("Exception replaying {}: {}", orderedEvent, rte) - } - ) + if (parallel) { + siriusLog.parallelForeach(bootstrapEvent) + } else { + siriusLog.foreach(bootstrapEvent) + } requestHandler.onBootstrapComplete() val totalBootstrapTime = System.currentTimeMillis - start bootstrapTime = Some(totalBootstrapTime) @@ -142,6 +135,18 @@ class StateSup(requestHandler: RequestHandler, } } + private def bootstrapEvent(orderedEvent : OrderedEvent): Unit = + try { + orderedEvent.request match { + case Put(key, body) => requestHandler.handlePut(orderedEvent.sequence, key, body) + case Delete(key) => requestHandler.handleDelete(orderedEvent.sequence, key) + } + } catch { + case rte: RuntimeException => + eventReplayFailureCount += 1 + logger.error("Exception replaying {}: {}", orderedEvent, rte) + } + trait StateInfoMBean { def getEventReplayFailureCount: Long def getBootstrapTime: String diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberPair.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberPair.scala index 439549d1..b8e9d691 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberPair.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberPair.scala @@ -30,9 +30,9 @@ object UberPair { * * @return an instantiated UberStoreFilePair */ - def apply(baseDir: String, startingSeq: Long, fileHandleFactory: UberDataFileHandleFactory): UberPair = { + def apply(baseDir: String, startingSeq: Long, fileHandleFactory: UberDataFileHandleFactory, validateChecksum: Boolean): UberPair = { val baseName = "%s/%s".format(baseDir, startingSeq) - val dataFile = UberDataFile("%s.data".format(baseName), fileHandleFactory) + val dataFile = UberDataFile("%s.data".format(baseName), fileHandleFactory, validateChecksum) val index = DiskOnlySeqIndex("%s.index".format(baseName)) repairIndex(index, dataFile) new UberPair(dataFile, index) diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberStore.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberStore.scala index 10e026cf..5b9c0b23 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberStore.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/UberStore.scala @@ -39,8 +39,10 @@ object UberStore { } val fileHandleFactory = UberDataFileHandleFactory(siriusConfig) + val skipChecksumValidation = siriusConfig.getProp(SiriusConfiguration.LOG_SKIP_CHECKSUM_VALIDATION, false) + val validateChecksum = !skipChecksumValidation - new UberStore(baseDir, UberPair(baseDir, 1L, fileHandleFactory)) + new UberStore(baseDir, UberPair(baseDir, 1L, fileHandleFactory, validateChecksum)) } /** diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/common/Checksummer.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/common/Checksummer.scala index e464623a..d6e63392 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/common/Checksummer.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/common/Checksummer.scala @@ -19,6 +19,16 @@ package com.comcast.xfinity.sirius.uberstore.common * Trait supplying checksumming capabilities */ trait Checksummer { + /** + * Determines if the array of bytes has a calculated checksum + * that matches the provided checksum + * + * @param chksum the checksum + * @param bytes Array[Byte] to checksum + * @return if the checksum of the bytes matches the provided checksum + */ + def validate(bytes: Array[Byte], chksum: Long): Boolean = + checksum(bytes) == chksum /** * Given an array of bytes will calculate a Long checksum diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/common/SkipValidationChecksummer.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/common/SkipValidationChecksummer.scala new file mode 100644 index 00000000..bcb8974d --- /dev/null +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/common/SkipValidationChecksummer.scala @@ -0,0 +1,8 @@ +package com.comcast.xfinity.sirius.uberstore.common + +trait SkipValidationChecksummer extends Checksummer { + /** + * Skips calculating and validating the checksum on the Array[Byte] + */ + override def validate(bytes: Array[Byte], chksum: Long): Boolean = true +} diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberDataFile.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberDataFile.scala index f77bfaa6..42526b80 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberDataFile.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberDataFile.scala @@ -16,7 +16,7 @@ package com.comcast.xfinity.sirius.uberstore.data import com.comcast.xfinity.sirius.api.impl.OrderedEvent -import com.comcast.xfinity.sirius.uberstore.common.Fnv1aChecksummer +import com.comcast.xfinity.sirius.uberstore.common.{Fnv1aChecksummer, SkipValidationChecksummer} import scala.annotation.tailrec @@ -35,8 +35,11 @@ object UberDataFile { * * @return fully constructed UberDataFile */ - def apply(dataFileName: String, fileHandleFactory: UberDataFileHandleFactory): UberDataFile = { - val fileOps = new UberStoreBinaryFileOps with Fnv1aChecksummer + def apply(dataFileName: String, fileHandleFactory: UberDataFileHandleFactory, validateChecksum: Boolean): UberDataFile = { + val fileOps = if (validateChecksum) + new UberStoreBinaryFileOps with Fnv1aChecksummer + else + new UberStoreBinaryFileOps with Fnv1aChecksummer with SkipValidationChecksummer val codec = new BinaryEventCodec new UberDataFile(dataFileName, fileHandleFactory, fileOps, codec) } diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberStoreBinaryFileOps.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberStoreBinaryFileOps.scala index b8819c73..575fa41e 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberStoreBinaryFileOps.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/data/UberStoreBinaryFileOps.scala @@ -53,7 +53,7 @@ class UberStoreBinaryFileOps extends UberStoreFileOps { } else { val (bodyLen, chksum) = readHeader(readHandle) val body = readBody(readHandle, bodyLen) - if (chksum == checksum(body)) { + if (validate(body, chksum)) { Some(body) // [that i used to know | to love] } else { throw new IllegalStateException("File corrupted at offset " + readHandle.offset()) diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/Segment.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/Segment.scala index d7497e33..a5d8f6f0 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/Segment.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/Segment.scala @@ -22,18 +22,6 @@ import java.io.File object Segment { - /** - * Create an Segment based in baseDir named "name". - * - * @param base directory containing the Segment - * @param name the name of this dir - * - * @return an Segment instance, fully repaired and usable - */ - def apply(base: File, name: String, fileHandleFactory: UberDataFileHandleFactory): Segment = { - apply(new File(base, name), fileHandleFactory) - } - /** * Create an Segment at the specified location. * @@ -42,7 +30,7 @@ object Segment { * * @return an Segment instance, fully repaired and usable */ - def apply(location: File, fileHandleFactory: UberDataFileHandleFactory): Segment = { + def apply(location: File, fileHandleFactory: UberDataFileHandleFactory, validateChecksum: Boolean = true): Segment = { location.mkdirs() val dataFile = new File(location, "data") @@ -50,7 +38,7 @@ object Segment { val compactionFlagFile = new File(location, "keys-collected") val internalCompactionFlagFile = new File(location, "internally-compacted") - val data = UberDataFile(dataFile.getAbsolutePath, fileHandleFactory) + val data = UberDataFile(dataFile.getAbsolutePath, fileHandleFactory, validateChecksum) val index = DiskOnlySeqIndex(indexFile.getAbsolutePath) val compactionFlag = FlagFile(compactionFlagFile.getAbsolutePath) val internalCompactionFlag = FlagFile(internalCompactionFlagFile.getAbsolutePath) diff --git a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala index 5be76948..4b8b6fe7 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStore.scala @@ -16,7 +16,6 @@ package com.comcast.xfinity.sirius.uberstore.segmented import java.io.{File => JFile} - import better.files.File import com.comcast.xfinity.sirius.api.SiriusConfiguration import com.comcast.xfinity.sirius.api.impl.OrderedEvent @@ -87,8 +86,10 @@ object SegmentedUberStore { val MAX_EVENTS_PER_SEGMENT = siriusConfig.getProp(SiriusConfiguration.LOG_EVENTS_PER_SEGMENT, 1000000L) val fileHandleFactory = UberDataFileHandleFactory(siriusConfig) + val skipChecksumValidation = siriusConfig.getProp(SiriusConfiguration.LOG_SKIP_CHECKSUM_VALIDATION, false) + val validateChecksum = !skipChecksumValidation - def buildSegment(location: JFile) = Segment(location, fileHandleFactory) + def buildSegment(location: JFile) = Segment(location, fileHandleFactory, validateChecksum) val segmentedCompactor = SegmentedCompactor(siriusConfig, buildSegment) new SegmentedUberStore(new JFile(base), MAX_EVENTS_PER_SEGMENT, segmentedCompactor, buildSegment) @@ -136,6 +137,11 @@ class SegmentedUberStore private[segmented] (base: JFile, */ def getNextSeq = nextSeq + override def parallelForeach[T](fun: OrderedEvent => T): Unit = { + ParallelHelpers.parallelize(readOnlyDirs :+ liveDir) + .foreach(_.foldLeftRange(0, Long.MaxValue)(())((_, e) => fun(e))) + } + /** * @inheritdoc */ diff --git a/src/main/scala/com/comcast/xfinity/sirius/writeaheadlog/SiriusLog.scala b/src/main/scala/com/comcast/xfinity/sirius/writeaheadlog/SiriusLog.scala index b8b12070..53254349 100644 --- a/src/main/scala/com/comcast/xfinity/sirius/writeaheadlog/SiriusLog.scala +++ b/src/main/scala/com/comcast/xfinity/sirius/writeaheadlog/SiriusLog.scala @@ -36,6 +36,13 @@ trait SiriusLog { */ def foreach[T](fun: OrderedEvent => T): Unit = foldLeft(())((_, e) => fun(e)) + /** + * Apply fun to each entry in the log in parallel and potentially out of order + * + * @param fun function to apply + */ + def parallelForeach[T](fun: OrderedEvent => T): Unit = foreach[T](fun) + /** * Fold left across the log entries * @param acc0 initial accumulator value diff --git a/src/test/scala/com/comcast/xfinity/sirius/uberstore/common/Fnv1aChecksummerTest.scala b/src/test/scala/com/comcast/xfinity/sirius/uberstore/common/Fnv1aChecksummerTest.scala index 34fbfc86..193beb72 100644 --- a/src/test/scala/com/comcast/xfinity/sirius/uberstore/common/Fnv1aChecksummerTest.scala +++ b/src/test/scala/com/comcast/xfinity/sirius/uberstore/common/Fnv1aChecksummerTest.scala @@ -32,4 +32,9 @@ class Fnv1aChecksummerTest extends NiceTest { val referenceBytes = "http://en.wikipedia.org/wiki/Fowler_Noll_Vo_hash".getBytes assert(-2758076559093427003L === underTest.checksum(referenceBytes)) } + + it ("validates the checksum") { + val referenceBytes = "http://en.wikipedia.org/wiki/Fowler_Noll_Vo_hash".getBytes + assert(true === underTest.validate(referenceBytes, -2758076559093427003L)) + } } diff --git a/src/test/scala/com/comcast/xfinity/sirius/uberstore/common/SkipValidationChecksummerTest.scala b/src/test/scala/com/comcast/xfinity/sirius/uberstore/common/SkipValidationChecksummerTest.scala new file mode 100644 index 00000000..e6341e84 --- /dev/null +++ b/src/test/scala/com/comcast/xfinity/sirius/uberstore/common/SkipValidationChecksummerTest.scala @@ -0,0 +1,19 @@ +package com.comcast.xfinity.sirius.uberstore.common + +import com.comcast.xfinity.sirius.NiceTest + +class SkipValidationChecksummerTest extends NiceTest { + val underTest = new Object with Fnv1aChecksummer with SkipValidationChecksummer + + it ("must match the reference impl") { + // from: + // http://trac.tools.ietf.org/wg/tls/draft-ietf-tls-cached-info/draft-ietf-tls-cached-info-06-from-05.diff.txt + val referenceBytes = "http://en.wikipedia.org/wiki/Fowler_Noll_Vo_hash".getBytes + assert(-2758076559093427003L === underTest.checksum(referenceBytes)) + } + + it ("skips validating the checksum") { + val referenceBytes = "http://en.wikipedia.org/wiki/Fowler_Noll_Vo_hash".getBytes + assert(true === underTest.validate(referenceBytes, 0L)) + } +} diff --git a/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentTest.scala b/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentTest.scala index 9b9c92bd..5001e9d8 100644 --- a/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentTest.scala +++ b/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentTest.scala @@ -58,7 +58,7 @@ class SegmentTest extends NiceTest with BeforeAndAfterAll { val fileHandleFactory: UberDataFileHandleFactory = RandomAccessFileHandleFactory def buildSegment(base: JFile, name: String) = - Segment(base, name, fileHandleFactory) + Segment(new JFile(base, name), fileHandleFactory) override def afterAll(): Unit = { File(tempDir.getPath).delete() diff --git a/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedCompactorTest.scala b/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedCompactorTest.scala index 78177724..807a741f 100644 --- a/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedCompactorTest.scala +++ b/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedCompactorTest.scala @@ -51,7 +51,7 @@ class SegmentedCompactorTest extends NiceTest with BeforeAndAfterAll { } def buildSegment(base: JFile, name: String): Segment = { - Segment(base, name, fileHandleFactory) + Segment(new JFile(base, name), fileHandleFactory) } def buildSegment(fullPath: String): Segment = { diff --git a/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStoreTest.scala b/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStoreTest.scala index acee7ee5..9fe9bb5f 100644 --- a/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStoreTest.scala +++ b/src/test/scala/com/comcast/xfinity/sirius/uberstore/segmented/SegmentedUberStoreTest.scala @@ -16,6 +16,8 @@ package com.comcast.xfinity.sirius.uberstore.segmented +import scala.collection.concurrent._ + import com.comcast.xfinity.sirius.NiceTest import java.io.{File => JFile} @@ -77,7 +79,7 @@ class SegmentedUberStoreTest extends NiceTest { } def buildSegment(base: JFile, name: String): Segment = - Segment(base, name, fileHandleFactory) + Segment(new JFile(base, name), fileHandleFactory) def buildSegment(location: JFile): Segment = Segment(location, fileHandleFactory) @@ -195,6 +197,21 @@ class SegmentedUberStoreTest extends NiceTest { } } + describe("parallelForeach") { + it("should bootstrap the uberstore in parallel") { + createPopulatedSegment(dir, "1", Range.inclusive(1, 3).toList, isApplied = true) + createPopulatedSegment(dir, "2", Range.inclusive(4, 6).toList, isApplied = true) + createPopulatedSegment(dir, "3", Range.inclusive(7, 9).toList, isApplied = true) + val config = new SiriusConfiguration + config.setProp(SiriusConfiguration.LOG_PARALLEL_ENABLED, true) + uberstore = SegmentedUberStore(dir.getAbsolutePath, config) + val map = new TrieMap[Long, SiriusRequest]() + uberstore.parallelForeach(event => map.put(event.sequence, event.request)) + + assert(map.size == 9) + } + } + describe("close") { it("should close all of the associated uberdirs") { uberstore.close()