Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -52,35 +52,35 @@ class BlockGeneratorSuite extends SparkFunSuite with BeforeAndAfter with TimeLim
val clock = new ManualClock()

require(blockIntervalMs > 5)
require(listener.onAddDataCalled === false)
require(listener.onGenerateBlockCalled === false)
require(listener.onPushBlockCalled === false)
require(!listener.onAddDataCalled)
require(!listener.onGenerateBlockCalled)
require(!listener.onPushBlockCalled)

// Verify that creating the generator does not start it
blockGenerator = new BlockGenerator(listener, 0, conf, clock)
assert(blockGenerator.isActive() === false, "block generator active before start()")
assert(blockGenerator.isStopped() === false, "block generator stopped before start()")
assert(listener.onAddDataCalled === false)
assert(listener.onGenerateBlockCalled === false)
assert(listener.onPushBlockCalled === false)
assert(!blockGenerator.isActive(), "block generator active before start()")
assert(!blockGenerator.isStopped(), "block generator stopped before start()")
assert(!listener.onAddDataCalled)
assert(!listener.onGenerateBlockCalled)
assert(!listener.onPushBlockCalled)

// Verify start marks the generator active, but does not call the callbacks
blockGenerator.start()
assert(blockGenerator.isActive(), "block generator active after start()")
assert(blockGenerator.isStopped() === false, "block generator stopped after start()")
assert(!blockGenerator.isStopped(), "block generator stopped after start()")
withClue("callbacks called before adding data") {
assert(listener.onAddDataCalled === false)
assert(listener.onGenerateBlockCalled === false)
assert(listener.onPushBlockCalled === false)
assert(!listener.onAddDataCalled)
assert(!listener.onGenerateBlockCalled)
assert(!listener.onPushBlockCalled)
}

// Verify whether addData() adds data that is present in generated blocks
val data1 = 1 to 10
data1.foreach { blockGenerator.addData _ }
withClue("callbacks called on adding data without metadata and without block generation") {
assert(listener.onAddDataCalled === false) // should be called only with addDataWithCallback()
assert(listener.onGenerateBlockCalled === false)
assert(listener.onPushBlockCalled === false)
assert(!listener.onAddDataCalled) // should be called only with addDataWithCallback()
assert(!listener.onGenerateBlockCalled)
assert(!listener.onPushBlockCalled)
}
clock.advance(blockIntervalMs) // advance clock to generate blocks
withClue("blocks not generated or pushed") {
Expand All @@ -90,7 +90,7 @@ class BlockGeneratorSuite extends SparkFunSuite with BeforeAndAfter with TimeLim
}
}
listener.pushedData.asScala.toSeq should contain theSameElementsInOrderAs (data1)
assert(listener.onAddDataCalled === false) // should be called only with addDataWithCallback()
assert(!listener.onAddDataCalled) // should be called only with addDataWithCallback()

// Verify addDataWithCallback() add data+metadata and callbacks are called correctly
val data2 = 11 to 20
Expand Down Expand Up @@ -146,10 +146,10 @@ class BlockGeneratorSuite extends SparkFunSuite with BeforeAndAfter with TimeLim
val listener = new TestBlockGeneratorListener
val clock = new ManualClock()
blockGenerator = new BlockGenerator(listener, 0, conf, clock)
require(listener.onGenerateBlockCalled === false)
require(!listener.onGenerateBlockCalled)
blockGenerator.start()
assert(blockGenerator.isActive(), "block generator")
assert(blockGenerator.isStopped() === false)
assert(!blockGenerator.isStopped())

val data = 1 to 1000
data.foreach { blockGenerator.addData _ }
Expand All @@ -161,9 +161,9 @@ class BlockGeneratorSuite extends SparkFunSuite with BeforeAndAfter with TimeLim
clock.advance(1) // to make sure that the timer for another interval to complete
val thread = stopBlockGenerator(blockGenerator)
eventually(timeout(1.second), interval(10.milliseconds)) {
assert(blockGenerator.isActive() === false)
assert(!blockGenerator.isActive())
}
assert(blockGenerator.isStopped() === false)
assert(!blockGenerator.isStopped())

// Verify that data cannot be added
intercept[SparkException] {
Expand Down Expand Up @@ -192,11 +192,11 @@ class BlockGeneratorSuite extends SparkFunSuite with BeforeAndAfter with TimeLim

// Verify that the final data is present in the final generated block and
// pushed before complete stop
assert(blockGenerator.isStopped() === false) // generator has not stopped yet
assert(!blockGenerator.isStopped()) // generator has not stopped yet
eventually(timeout(10.seconds), interval(10.milliseconds)) {
// Keep calling `advance` to avoid blocking forever in `clock.waitTillTime`
clock.advance(blockIntervalMs)
assert(thread.isAlive === false)
assert(!thread.isAlive)
}
assert(blockGenerator.isStopped()) // generator has finally been completely stopped
assert(listener.pushedData.asScala.toSeq === data, "All data not pushed by stop()")
Expand All @@ -211,7 +211,7 @@ class BlockGeneratorSuite extends SparkFunSuite with BeforeAndAfter with TimeLim
}
blockGenerator = new BlockGenerator(listener, 0, conf)
blockGenerator.start()
assert(listener.onErrorCalled === false)
assert(!listener.onErrorCalled)
blockGenerator.addData(1)
eventually(timeout(1.second), interval(10.milliseconds)) {
assert(listener.onErrorCalled)
Expand Down