diff --git a/core/amber/src/main/scala/edu/uci/ics/amber/AmberApp.scala b/core/amber/src/main/scala/edu/uci/ics/amber/AmberApp.scala index ec77d297c88..88ac22e51ee 100644 --- a/core/amber/src/main/scala/edu/uci/ics/amber/AmberApp.scala +++ b/core/amber/src/main/scala/edu/uci/ics/amber/AmberApp.scala @@ -6,7 +6,6 @@ import akka.actor.{ActorRef, ActorSystem, Props} import akka.util.Timeout import com.typesafe.config.{Config, ConfigFactory} import edu.uci.ics.amber.clustering.ClusterListener -import edu.uci.ics.amber.engine.architecture.breakpoint.globalbreakpoint.ConditionalGlobalBreakpoint import edu.uci.ics.amber.engine.architecture.controller.Controller import edu.uci.ics.amber.engine.common.Constants import edu.uci.ics.amber.engine.common.ambermessage.ControlMessage.{Pause, Resume, Start} @@ -256,15 +255,15 @@ object AmberApp { //if (countbp.isDefined && current == 2) { // controller ! PassBreakpointTo("Filter", new CountGlobalBreakpoint("CountBreakpoint", countbp.get)) //} - if (conditionalbp.isDefined) { - controller ! PassBreakpointTo( - "KeywordSearch", - new ConditionalGlobalBreakpoint( - "ConditionalBreakpoint", - x => x.getString(15).contains(conditionalbp) - ) - ) - } +// if (conditionalbp.isDefined) { +// controller ! PassBreakpointTo( +// "KeywordSearch", +// new ConditionalGlobalBreakpoint( +// "ConditionalBreakpoint", +// x => x.getString(15).contains(conditionalbp) +// ) +// ) +// } controller ! Start println("workflow started!") case "pause" => diff --git a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/globalbreakpoint/ConditionalGlobalBreakpoint.scala b/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/globalbreakpoint/ConditionalGlobalBreakpoint.scala deleted file mode 100644 index fa90aa31a10..00000000000 --- a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/globalbreakpoint/ConditionalGlobalBreakpoint.scala +++ /dev/null @@ -1,73 +0,0 @@ -package edu.uci.ics.amber.engine.architecture.breakpoint.globalbreakpoint - -import edu.uci.ics.amber.engine.architecture.breakpoint.FaultedTuple -import edu.uci.ics.amber.engine.architecture.breakpoint.localbreakpoint.{ - ConditionalBreakpoint, - CountBreakpoint, - LocalBreakpoint -} -import edu.uci.ics.amber.engine.common.AdvancedMessageSending -import edu.uci.ics.amber.engine.common.ambermessage.WorkerMessage.{ - AssignBreakpoint, - QueryBreakpoint, - RemoveBreakpoint -} -import edu.uci.ics.amber.engine.common.tuple.ITuple -import akka.actor.ActorRef -import akka.event.LoggingAdapter -import akka.util.Timeout - -import scala.collection.mutable -import scala.collection.mutable.ArrayBuffer -import scala.concurrent.ExecutionContext - -class ConditionalGlobalBreakpoint(id: String, val predicate: ITuple => Boolean) - extends GlobalBreakpoint(id) { - - var localbreakpoints: ArrayBuffer[(ActorRef, LocalBreakpoint)] = - new ArrayBuffer[(ActorRef, LocalBreakpoint)]() - - override def acceptImpl(sender: ActorRef, localBreakpoint: LocalBreakpoint): Unit = { - if (localBreakpoint.isTriggered) { - localbreakpoints.append((sender, localBreakpoint)) - } - } - - override def isTriggered: Boolean = localbreakpoints.nonEmpty - - override def partitionImpl(layer: Array[ActorRef])(implicit - timeout: Timeout, - ec: ExecutionContext, - id: String, - version: Long - ): Iterable[ActorRef] = { - for (x <- layer) { - AdvancedMessageSending.blockingAskWithRetry( - x, - AssignBreakpoint(new ConditionalBreakpoint(predicate)), - 10 - ) - } - layer - } - - override def report(map: mutable.HashMap[(ActorRef, FaultedTuple), ArrayBuffer[String]]): Unit = { - for (i <- localbreakpoints) { - val k = (i._1, new FaultedTuple(i._2.triggeredTuple, i._2.triggeredTupleId, false)) - if (map.contains(k)) { - map(k).append("condition unsatisfied") - } else { - map(k) = ArrayBuffer[String]("condition unsatisfied") - } - } - localbreakpoints.clear() - } - - override def isCompleted: Boolean = false - - override def reset(): Unit = { - super.reset() - localbreakpoints = new ArrayBuffer[(ActorRef, LocalBreakpoint)]() - } - -} diff --git a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/globalbreakpoint/CountGlobalBreakpoint.scala b/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/globalbreakpoint/CountGlobalBreakpoint.scala deleted file mode 100644 index d061cdf1a5b..00000000000 --- a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/globalbreakpoint/CountGlobalBreakpoint.scala +++ /dev/null @@ -1,93 +0,0 @@ -package edu.uci.ics.amber.engine.architecture.breakpoint.globalbreakpoint -import edu.uci.ics.amber.engine.architecture.breakpoint.FaultedTuple -import edu.uci.ics.amber.engine.architecture.breakpoint.localbreakpoint.{ - CountBreakpoint, - LocalBreakpoint -} -import edu.uci.ics.amber.engine.common.AdvancedMessageSending -import edu.uci.ics.amber.engine.common.ambermessage.WorkerMessage.{ - AssignBreakpoint, - QueryBreakpoint, - QueryTriggeredBreakpoints, - RemoveBreakpoint -} -import akka.actor.ActorRef -import akka.event.LoggingAdapter -import akka.util.Timeout - -import scala.collection.mutable -import scala.collection.mutable.ArrayBuffer -import scala.concurrent.ExecutionContext - -class CountGlobalBreakpoint(id: String, val target: Long) extends GlobalBreakpoint(id) { - - var current: Long = 0 - var localbreakpoints: ArrayBuffer[(ActorRef, LocalBreakpoint)] = - new ArrayBuffer[(ActorRef, LocalBreakpoint)]() - - override def acceptImpl(sender: ActorRef, localBreakpoint: LocalBreakpoint): Unit = { - current += localBreakpoint.asInstanceOf[CountBreakpoint].current - if (localBreakpoint.isTriggered) - localbreakpoints.append((sender, localBreakpoint)) - } - - override def isTriggered: Boolean = current == target - - override def partitionImpl(layer: Array[ActorRef])(implicit - timeout: Timeout, - ec: ExecutionContext, - id: String, - version: Long - ): Iterable[ActorRef] = { - val remaining = target - current - var currentSum = 0L - val length = layer.length - var i = 0 - if (remaining / length > 0) { - while (i < length - 1) { - AdvancedMessageSending.blockingAskWithRetry( - layer(i), - AssignBreakpoint(new CountBreakpoint(remaining / length)), - 10 - ) - currentSum += remaining / length - i += 1 - } - AdvancedMessageSending.blockingAskWithRetry( - layer.last, - AssignBreakpoint(new CountBreakpoint(remaining - currentSum)), - 10 - ) - layer - } else { - AdvancedMessageSending.blockingAskWithRetry( - layer.last, - AssignBreakpoint(new CountBreakpoint(remaining)), - 10 - ) - Array(layer.last) - } - } - - override def isRepartitionRequired: Boolean = unReportedWorkers.isEmpty && target != current - - override def report(map: mutable.HashMap[(ActorRef, FaultedTuple), ArrayBuffer[String]]): Unit = { - for (i <- localbreakpoints) { - val k = (i._1, new FaultedTuple(i._2.triggeredTuple, i._2.triggeredTupleId, false)) - if (map.contains(k)) { - map(k).append(s"count reached $target") - } else { - map(k) = ArrayBuffer[String](s"count reached $target") - } - } - localbreakpoints.clear() - } - - override def isCompleted: Boolean = isTriggered - - override def reset(): Unit = { - super.reset() - current = 0 - } - -} diff --git a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/globalbreakpoint/ExceptionGlobalBreakpoint.scala b/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/globalbreakpoint/ExceptionGlobalBreakpoint.scala deleted file mode 100644 index 09ab386ab7c..00000000000 --- a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/globalbreakpoint/ExceptionGlobalBreakpoint.scala +++ /dev/null @@ -1,66 +0,0 @@ -package edu.uci.ics.amber.engine.architecture.breakpoint.globalbreakpoint - -import edu.uci.ics.amber.engine.architecture.breakpoint.FaultedTuple -import edu.uci.ics.amber.engine.architecture.breakpoint.localbreakpoint.{ - ConditionalBreakpoint, - ExceptionBreakpoint, - LocalBreakpoint -} -import edu.uci.ics.amber.engine.common.AdvancedMessageSending -import edu.uci.ics.amber.engine.common.ambermessage.WorkerMessage.AssignBreakpoint -import edu.uci.ics.amber.engine.common.tuple.ITuple -import akka.actor.ActorRef -import akka.event.LoggingAdapter -import akka.util.Timeout - -import scala.collection.mutable -import scala.collection.mutable.ArrayBuffer -import scala.concurrent.ExecutionContext - -class ExceptionGlobalBreakpoint(id: String) extends GlobalBreakpoint(id) { - var exceptions: ArrayBuffer[(ActorRef, ExceptionBreakpoint)] = - new ArrayBuffer[(ActorRef, ExceptionBreakpoint)]() - - override def acceptImpl(sender: ActorRef, localBreakpoint: LocalBreakpoint): Unit = { - if (localBreakpoint.isTriggered) { - exceptions.append((sender, localBreakpoint.asInstanceOf[ExceptionBreakpoint])) - } - } - - override def isTriggered: Boolean = exceptions.nonEmpty - - override def partitionImpl(layer: Array[ActorRef])(implicit - timeout: Timeout, - ec: ExecutionContext, - id: String, - version: Long - ): Iterable[ActorRef] = { - for (x <- layer) { - AdvancedMessageSending.blockingAskWithRetry( - x, - AssignBreakpoint(new ExceptionBreakpoint()), - 10 - ) - } - layer - } - - override def report(map: mutable.HashMap[(ActorRef, FaultedTuple), ArrayBuffer[String]]): Unit = { - for (i <- exceptions) { - val k = (i._1, new FaultedTuple(i._2.triggeredTuple, i._2.triggeredTupleId, i._2.isInput)) - if (map.contains(k)) { - map(k).append(i._2.error.toString) - } else { - map(k) = ArrayBuffer[String](i._2.error.toString) - } - } - exceptions.clear() - } - - override def isCompleted: Boolean = false - - override def reset(): Unit = { - super.reset() - exceptions = new ArrayBuffer[(ActorRef, ExceptionBreakpoint)]() - } -} diff --git a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/localbreakpoint/ConditionalBreakpoint.scala b/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/localbreakpoint/ConditionalBreakpoint.scala deleted file mode 100644 index 528dea50540..00000000000 --- a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/localbreakpoint/ConditionalBreakpoint.scala +++ /dev/null @@ -1,22 +0,0 @@ -package edu.uci.ics.amber.engine.architecture.breakpoint.localbreakpoint - -import edu.uci.ics.amber.engine.common.tuple.ITuple - -class ConditionalBreakpoint(val predicate: ITuple => Boolean)(implicit id: String, version: Long) - extends LocalBreakpoint(id, version) { - - var _isTriggered = false - - override def accept(tuple: ITuple): Unit = { - _isTriggered = predicate(tuple) - } - - override def isTriggered: Boolean = _isTriggered - - override def isDirty: Boolean = isTriggered - - override def reset(): Unit = { - super.reset() - _isTriggered = false - } -} diff --git a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/localbreakpoint/CountBreakpoint.scala b/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/localbreakpoint/CountBreakpoint.scala deleted file mode 100644 index 71be3a1c41d..00000000000 --- a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/localbreakpoint/CountBreakpoint.scala +++ /dev/null @@ -1,22 +0,0 @@ -package edu.uci.ics.amber.engine.architecture.breakpoint.localbreakpoint - -import edu.uci.ics.amber.engine.common.tuple.ITuple - -class CountBreakpoint(val target: Long)(implicit id: String, version: Long) - extends LocalBreakpoint(id, version) { - - var current: Long = 0 - - override def accept(tuple: ITuple): Unit = { - current += 1 - } - - override def isTriggered: Boolean = current == target - - override def isDirty: Boolean = isReported - - override def reset(): Unit = { - super.reset() - current = 0 - } -} diff --git a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/localbreakpoint/ExceptionBreakpoint.scala b/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/localbreakpoint/ExceptionBreakpoint.scala deleted file mode 100644 index 2cb4504b650..00000000000 --- a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/breakpoint/localbreakpoint/ExceptionBreakpoint.scala +++ /dev/null @@ -1,19 +0,0 @@ -package edu.uci.ics.amber.engine.architecture.breakpoint.localbreakpoint -import edu.uci.ics.amber.engine.common.tuple.ITuple - -class ExceptionBreakpoint()(implicit id: String, version: Long) - extends LocalBreakpoint(id, version) { - var error: Exception = _ - override def accept(tuple: ITuple): Unit = { - //empty - } - - override def isTriggered: Boolean = error != null - - override def isDirty: Boolean = isTriggered - - override def reset(): Unit = { - super.reset() - error = null - } -} diff --git a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/controller/Controller.scala b/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/controller/Controller.scala index 1903056b1ac..a86f8afe548 100644 --- a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/controller/Controller.scala +++ b/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/controller/Controller.scala @@ -1,10 +1,7 @@ package edu.uci.ics.amber.engine.architecture.controller import edu.uci.ics.amber.clustering.ClusterListener.GetAvailableNodeAddresses -import edu.uci.ics.amber.engine.architecture.breakpoint.globalbreakpoint.{ - ExceptionGlobalBreakpoint, - GlobalBreakpoint -} +import edu.uci.ics.amber.engine.architecture.breakpoint.globalbreakpoint.GlobalBreakpoint import edu.uci.ics.amber.engine.architecture.breakpoint.globalbreakpoint.GlobalBreakpoint import edu.uci.ics.amber.engine.architecture.controller.ControllerEvent.{ BreakpointTriggered, @@ -29,11 +26,6 @@ import edu.uci.ics.amber.engine.architecture.linksemantics.{ LocalPartialToOne, OperatorLink } -import edu.uci.ics.amber.engine.architecture.principal.{ - Principal, - PrincipalState, - PrincipalStatistics -} import edu.uci.ics.amber.engine.common.amberexception.WorkflowRuntimeException import edu.uci.ics.amber.engine.common.ambermessage.ControllerMessage._ import edu.uci.ics.amber.engine.common.ambermessage.ControlMessage._ @@ -110,6 +102,7 @@ import edu.uci.ics.amber.engine.common.ambermessage.WorkerMessage.{ import edu.uci.ics.amber.error.WorkflowRuntimeError import edu.uci.ics.amber.engine.architecture.messaginglayer.NetworkCommunicationActor import edu.uci.ics.amber.engine.architecture.messaginglayer.NetworkCommunicationActor.RegisterActorRef +import edu.uci.ics.amber.engine.architecture.principal.{PrincipalState, PrincipalStatistics} import edu.uci.ics.amber.engine.common.ambertag.neo.VirtualIdentity import edu.uci.ics.amber.engine.common.rpc.AsyncRPCHandlerInitializer @@ -176,12 +169,6 @@ class Controller( val pauseTimer = Stopwatch.createUnstarted(); var recoveryMode = false -// def allPrincipals: Iterable[ActorRef] = principalStates.keys -// def unCompletedPrincipals: Iterable[ActorRef] = -// principalStates.filter(x => x._2 != PrincipalState.Completed).keys -// def allUnCompletedPrincipalStates: Iterable[PrincipalState.Value] = -// principalStates.filter(x => x._2 != PrincipalState.Completed).values - def allUnCompletedOperatorStates: Iterable[PrincipalState.Value] = operatorStateMap.filter(x => x._2 != PrincipalState.Completed).values val tau: FiniteDuration = Constants.defaultTau @@ -204,8 +191,6 @@ class Controller( var operatorToGlobalBreakpoints = new mutable.HashMap[OperatorIdentifier, mutable.AnyRefMap[String, GlobalBreakpoint]]() var operatorToPeriodicallyAskHandle = new mutable.HashMap[OperatorIdentifier, Cancellable]() - var operatorToWorkersTriggeredBreakpoint = - new mutable.HashMap[OperatorIdentifier, Iterable[ActorRef]]() var operatorToLayerCompletedCounter = new mutable.HashMap[OperatorIdentifier, mutable.HashMap[LayerTag, Int]]() val operatorToTimer = @@ -232,54 +217,6 @@ class Controller( Await .result(context.actorSelection("/user/cluster-info") ? GetAvailableNodeAddresses, 5.seconds) .asInstanceOf[Array[Address]] -// def getPrincipalNode(nodes: Array[Address]): Address = -// self.path.address //nodes(util.Random.nextInt(nodes.length)) - - private def queryExecuteStatistics(): Unit = {} - - //if checkpoint activated: - private def insertCheckpoint(from: OpExecConfig, to: OpExecConfig): Unit = { - //insert checkpoint barrier between 2 operators and delete the link between them - val topology = from.topology - val hashFunc = to.getShuffleHashFunction(topology.layers.last.tag) - val layerTag = LayerTag(from.tag, "checkpoint") - val path: String = layerTag.getGlobalIdentity - val numWorkers = topology.layers.last.numWorkers - val scanGen: Int => ISourceOperatorExecutor = i => - new HDFSFolderScanSourceOperatorExecutor(Constants.remoteHDFSPath, path + "/" + i, '|', null) - val lastLayer = topology.layers.last - val materializerLayer = new WorkerLayer( - layerTag, - i => new HashBasedMaterializer(path, i, hashFunc, numWorkers), - numWorkers, - FollowPrevious(), - OneOnEach() - ) - topology.layers :+= materializerLayer - topology.links :+= new LocalPartialToOne( - lastLayer, - materializerLayer, - Constants.defaultBatchSize, - 0 - ) - val scanLayer = new WorkerLayer( - LayerTag(to.tag, "from_checkpoint"), - scanGen, - topology.layers.last.numWorkers, - FollowPrevious(), - OneOnEach() - ) - val firstLayer = to.topology.layers.head - to.topology.layers +:= scanLayer - to.topology.links +:= new HashBasedShuffle( - scanLayer, - firstLayer, - Constants.defaultBatchSize, - hashFunc, - 0 - ) - - } final def safeRemoveAskHandle(): Unit = { if (periodicallyAskHandle != null) { @@ -288,38 +225,6 @@ class Controller( } } - private def killAndRecoverStage(): Unit = { -// val futuresNoSinkScan = operatorInCurrentStage -// .filter(x => x.operator.contains("Sink") && x.operator.contains("Scan")) -// .map { x => -// operatorStateMap(x) = PrincipalState.Running -// operatorToWorkerLayers(x).foreach { layers => -// layers.layer(0) ! Reset( -// layers.getFirstMetadata, -// Seq(operatorToReceivedRecoveryInformation(x)(layers.tagForFirst)) -// ) -// operatorToWorkerStateMap(x)(layers.layer(0)) = WorkerState.Ready -// } -// } -// .asJava -// val tasksNoSinkScan = Futures.sequence(futuresNoSinkScan, ec) -// Await.result(tasksNoSinkScan, 5.minutes) -// Thread.sleep(2000) -// val futuresScan = principalInCurrentStage -// .filter(x => principalBiMap.inverse().get(x).operator.contains("Scan")) -// .map { x => -// principalStates(x) = PrincipalState.Running -// AdvancedMessageSending.nonBlockingAskWithRetry(x, KillAndRecover, 5, 0)(3.minutes, ec, log) -// } -// .asJava -// val tasksScan = Futures.sequence(futuresScan, ec) -// Await.result(tasksScan, 5.minutes) -// if (this.eventListener.recoveryStartedListener != null) { -// this.eventListener.recoveryStartedListener.apply() -// } -// context.become(pausing) - } - private def aggregateWorkerInputRowCount(opIdentifier: OperatorIdentifier): Long = { operatorToWorkerStatisticsMap(opIdentifier) .filter(e => operatorToWorkerLayers(opIdentifier).head.layer.contains(e._1)) @@ -468,19 +373,6 @@ class Controller( safeRemoveAskOperatorHandle(startOp) operatorToPeriodicallyAskHandle(startOp) = context.system.scheduler.schedule(30.seconds, 30.seconds, self, EnforceStateCheck(startOp)) - // context.become(initializing) - // unstashAll() - if (!recoveryMode) { - val breakpointToAssign = new ExceptionGlobalBreakpoint( - startOp.operator + "-ExceptionBreakpoint" - ) - operatorToGlobalBreakpoints(startOp)(breakpointToAssign.id) = breakpointToAssign - metadata.assignBreakpoint( - operatorToWorkerLayers(startOp), - operatorToWorkerStateMap(startOp), - breakpointToAssign - ) - } } private def initializingNextFrontier(): Unit = { @@ -533,7 +425,6 @@ class Controller( if (withCheckpoint) { for (n <- workflow.outLinks(k)) { if (workflow.operators(n).requiredShuffle) { - insertCheckpoint(workflow.operators(k), workflow.operators(n)) operatorsToWait.append(k) linksToIgnore.add((k, n)) } @@ -596,64 +487,6 @@ class Controller( val opId = workerToOperator(sender) if (setWorkerState(sender, state)) { state match { - case WorkerState.LocalBreakpointTriggered => - if ( - whenAllUncompletedWorkersBecome( - workerToOperator(sender), - WorkerState.LocalBreakpointTriggered - ) - ) { - //only one worker and it triggered breakpoint - safeRemoveAskOperatorHandle(workerToOperator(sender)) - operatorToPeriodicallyAskHandle(workerToOperator(sender)) = - context.system.scheduler.schedule( - 0.milliseconds, - 30.seconds, - self, - EnforceStateCheck(opId) - ) - operatorToWorkersTriggeredBreakpoint(workerToOperator(sender)) = - operatorToWorkerStateMap(workerToOperator(sender)).keys - operatorStateMap(workerToOperator(sender)) = PrincipalState.CollectingBreakpoints - matchNewOperatorStateAndTakeActionsInRunning(workerToOperator(sender)) - - } else { - //no tau involved since we know a very small tau works best - if (!operatorToStage2Timer(workerToOperator(sender)).isRunning) { - operatorToStage2Timer(workerToOperator(sender)).start() - } - if (operatorToStage1Timer(workerToOperator(sender)).isRunning) { - operatorToStage1Timer(workerToOperator(sender)).stop() - } - - operatorToWorkerStateMap(workerToOperator(sender)) - .filter(x => x._2 != WorkerState.Completed) - .keys - .foreach(worker => { - println( - s"Sending pause to ${worker.toString()} -- ${workerToOperator(sender).getGlobalIdentity}" - ) - }) - - val workersToPause: Iterable[ActorRef] = - operatorToWorkerStateMap(workerToOperator(sender)) - .filter(x => x._2 != WorkerState.Completed) - .keys - - context.system.scheduler - .scheduleOnce( - tau, - () => - workersToPause - .foreach(worker => worker ! Pause) - ) - - safeRemoveAskOperatorHandle(workerToOperator(sender)) - operatorStateMap(workerToOperator(sender)) = PrincipalState.Pausing - // context.become(pausing) - operatorToPeriodicallyAskHandle(workerToOperator(sender)) = context.system.scheduler - .schedule(30.seconds, 30.seconds, self, EnforceStateCheck(opId)) - } case WorkerState.Paused => if (areAllWorkersCompleted(workerToOperator(sender))) { safeRemoveAskOperatorHandle(workerToOperator(sender)) @@ -665,24 +498,6 @@ class Controller( safeRemoveAskOperatorHandle(workerToOperator(sender)) operatorStateMap(workerToOperator(sender)) = PrincipalState.Paused matchNewOperatorStateAndTakeActionsInRunning(workerToOperator(sender)) - } else if ( - unCompletedWorkerStates(workerToOperator(sender)) - .forall(x => x == WorkerState.Paused || x == WorkerState.LocalBreakpointTriggered) - ) { - operatorToWorkersTriggeredBreakpoint(workerToOperator(sender)) = - operatorToWorkerStateMap(workerToOperator(sender)) - .filter(_._2 == WorkerState.LocalBreakpointTriggered) - .keys - safeRemoveAskOperatorHandle(workerToOperator(sender)) - operatorToPeriodicallyAskHandle(workerToOperator(sender)) = context.system.scheduler - .schedule( - 1.milliseconds, - 30.seconds, - self, - EnforceStateCheck(opId) - ) - operatorStateMap(workerToOperator(sender)) = PrincipalState.CollectingBreakpoints - matchNewOperatorStateAndTakeActionsInRunning(workerToOperator(sender)) } case WorkerState.Completed => if (areAllWorkersCompleted(workerToOperator(sender))) { @@ -724,132 +539,10 @@ class Controller( operatorToReceivedTuples(workerToOperator(sender)).clear() operatorStateMap(workerToOperator(sender)) = PrincipalState.Paused matchNewOperatorStateAndTakeActionsInPausing(workerToOperator(sender)) - } else if ( - operatorToWorkerStateMap(workerToOperator(sender)) - .filter(x => x._2 != WorkerState.Completed) - .values - .forall(x => x == WorkerState.Paused || x == WorkerState.LocalBreakpointTriggered) - ) { - operatorToWorkersTriggeredBreakpoint(workerToOperator(sender)) = operatorToWorkerStateMap( - workerToOperator(sender) - ).filter(_._2 == WorkerState.LocalBreakpointTriggered).keys - safeRemoveAskOperatorHandle(workerToOperator(sender)) - operatorToPeriodicallyAskHandle(workerToOperator(sender)) = context.system.scheduler - .schedule( - 1.milliseconds, - 30.seconds, - self, - EnforceStateCheck(opId) - ) - operatorStateMap(workerToOperator(sender)) = PrincipalState.CollectingBreakpoints - matchNewOperatorStateAndTakeActionsInPausing(workerToOperator(sender)) - } - } - } - - private def handleWorkerStateReportsInCollBreakpoints(state: WorkerState.Value): Unit = { - controllerLogger.logInfo("collecting: " + sender + " to " + state) - val opId = workerToOperator(sender) - if (setWorkerState(sender, state)) { - if (unCompletedWorkerStates(workerToOperator(sender)).forall(_ == WorkerState.Paused)) { - //all breakpoint resolved, it's safe to report to controller and then Pause(on triggered, or user paused) else Resume - val map = new mutable.HashMap[(ActorRef, FaultedTuple), ArrayBuffer[String]] - for ( - i <- operatorToGlobalBreakpoints(workerToOperator(sender)).values.filter(_.isTriggered) - ) { - operatorToIsUserPaused(workerToOperator(sender)) = true //upgrade pause - i.report(map) - } - safeRemoveAskOperatorHandle(workerToOperator(sender)) - if (!operatorToIsUserPaused(workerToOperator(sender))) { - controllerLogger.logInfo("no global breakpoint triggered, continue") - operatorToIsUserPaused(workerToOperator(sender)) = false //reset - assert( - operatorToWorkerStateMap(workerToOperator(sender)) - .filter(x => x._2 != WorkerState.Completed) - .values - .nonEmpty - ) - operatorToWorkerStateMap(workerToOperator(sender)) - .filter(x => x._2 != WorkerState.Completed) - .keys - .foreach(worker => worker ! Resume) - safeRemoveAskOperatorHandle(workerToOperator(sender)) - operatorToPeriodicallyAskHandle(workerToOperator(sender)) = context.system.scheduler - .schedule(30.seconds, 30.seconds, self, EnforceStateCheck(opId)) - } else { - self ! Pause - context.parent ! ReportGlobalBreakpointTriggered( - map, - workflow.operators(workerToOperator(sender)).tag.operator - ) - if (this.eventListener.breakpointTriggeredListener != null) { - this.eventListener.breakpointTriggeredListener.apply( - BreakpointTriggered(map, workflow.operators(workerToOperator(sender)).tag.operator) - ) - } - controllerLogger.logInfo(map.toString()) - operatorStateMap(workerToOperator(sender)) = PrincipalState.Paused - controllerLogger.logInfo( - "user paused or global breakpoint triggered, pause. Stage1 cost = " + operatorToStage1Timer( - workerToOperator(sender) - ) - .toString() + " Stage2 cost =" + operatorToStage2Timer(workerToOperator(sender)) - .toString() - ) - } - if (operatorToStage2Timer(workerToOperator(sender)).isRunning) { - operatorToStage2Timer(workerToOperator(sender)).stop() - } - if (!operatorToStage1Timer(workerToOperator(sender)).isRunning) { - operatorToStage1Timer(workerToOperator(sender)).start() - } } } } - private[this] def handleBreakpointOnlyWorkerMessages: Receive = { - case ReportedTriggeredBreakpoints(bps) => - bps.foreach(x => { - val bp = operatorToGlobalBreakpoints(workerToOperator(sender))(x.id) - bp.accept(sender, x) - if (bp.needCollecting) { - //is not fully collected - bp.collect() - } else if (bp.isRepartitionRequired) { - //fully collected, but need repartition (e.g. count not reach target number) - //OR need Reset - workflow - .operators(workerToOperator(sender)) - .assignBreakpoint( - operatorToWorkerLayers(workerToOperator(sender)), - operatorToWorkerStateMap(workerToOperator(sender)), - bp - ) - } else if (bp.isCompleted) { - //fully collected and reach the target - bp.remove() - } - }) - case ReportedQueriedBreakpoint(bp) => - val gbp = operatorToGlobalBreakpoints(workerToOperator(sender))(bp.id) - if (gbp.accept(sender, bp) && !gbp.needCollecting) { - if (gbp.isRepartitionRequired) { - //fully collected, but need repartition (count not reach target number) - workflow - .operators(workerToOperator(sender)) - .assignBreakpoint( - operatorToWorkerLayers(workerToOperator(sender)), - operatorToWorkerStateMap(workerToOperator(sender)), - gbp - ) - } else if (gbp.isCompleted) { - //fully collected and reach the target - gbp.remove() - } - } - } - private def matchNewOperatorStateAndTakeActionsInRunning( opIdentifier: OperatorIdentifier ): Unit = { @@ -891,7 +584,6 @@ class Controller( self ! ContinuedInitialization } } - case PrincipalState.CollectingBreakpoints => case PrincipalState.Paused => if (operatorStateMap.values.forall(_ == PrincipalState.Completed)) { if (timer.isRunning) { @@ -1003,16 +695,6 @@ class Controller( if (whenAllUncompletedWorkersBecome(opIdentifier, WorkerState.Ready)) { safeRemoveAskOperatorHandle(opIdentifier) operatorToWorkerEdges(opIdentifier).foreach(x => x.link()) - operatorToGlobalBreakpoints(opIdentifier).values.foreach( - workflow - .operators(opIdentifier) - .assignBreakpoint( - operatorToWorkerLayers(opIdentifier), - operatorToWorkerStateMap(opIdentifier), - _ - ) - ) - operatorToWorkerStateMap(opIdentifier).keys.foreach(_ ! CheckRecovery) operatorStateMap(opIdentifier) = PrincipalState.Ready initializingNextFrontier() } @@ -1122,16 +804,7 @@ class Controller( case _ => //throw new AmberException("Invalid worker state received!") } case PassBreakpointTo(id: String, breakpoint: GlobalBreakpoint) => - val opTag = OperatorIdentifier(tag, id) - operatorToGlobalBreakpoints(opTag)(breakpoint.id) = breakpoint - controllerLogger.logInfo("assign breakpoint: " + breakpoint.id) - workflow - .operators(opTag) - .assignBreakpoint( - operatorToWorkerLayers(opTag), - operatorToWorkerStateMap(opTag), - breakpoint - ) + //TODO: invoke breakpoint assignment case msg => controllerLogger.logInfo("Stashing: " + msg) stash() @@ -1139,12 +812,9 @@ class Controller( } private[this] def running: Receive = { - disallowActorRefRelatedMessages orElse - handleBreakpointOnlyWorkerMessages orElse [Any, Unit] { + disallowActorRefRelatedMessages orElse [Any, Unit] { case LogErrorToFrontEnd(err: WorkflowRuntimeError) => controllerLogger.logError(err) - case KillAndRecover => - killAndRecoverStage() case QueryStatistics => operatorToWorkerLayers.keys.foreach(opIdentifier => { operatorToWorkerLayers(opIdentifier).foreach(l => { @@ -1158,10 +828,6 @@ class Controller( triggerStatusUpdateEvent(); case EnforceStateCheck(operatorIdentifier) => operatorStateMap(operatorIdentifier) match { - case PrincipalState.CollectingBreakpoints => - operatorToWorkersTriggeredBreakpoint(operatorIdentifier).foreach(x => - x ! QueryTriggeredBreakpoints - ) case PrincipalState.Pausing => for ((k, v) <- operatorToWorkerStateMap(operatorIdentifier)) { if (!allowedStatesOnPausing.contains(v)) { @@ -1177,8 +843,6 @@ class Controller( operatorStateMap(workerToOperator(sender)) = PrincipalState.Running } operatorStateMap(workerToOperator(sender)) match { - case PrincipalState.CollectingBreakpoints => - handleWorkerStateReportsInCollBreakpoints(state) case PrincipalState.Pausing => handleWorkerStateReportsInPausing(state) case _ => @@ -1239,8 +903,7 @@ class Controller( } private[this] def pausing: Receive = { - disallowActorRefRelatedMessages orElse - handleBreakpointOnlyWorkerMessages orElse [Any, Unit] { + disallowActorRefRelatedMessages orElse [Any, Unit] { case LogErrorToFrontEnd(err: WorkflowRuntimeError) => controllerLogger.logError(err) case QueryStatistics => @@ -1260,10 +923,6 @@ class Controller( ) case EnforceStateCheck(operatorIdentifier) => operatorStateMap(operatorIdentifier) match { - case PrincipalState.CollectingBreakpoints => - operatorToWorkersTriggeredBreakpoint(operatorIdentifier).foreach(x => - x ! QueryTriggeredBreakpoints - ) case _ => for ((k, v) <- operatorToWorkerStateMap(operatorIdentifier)) { if (!allowedStatesOnPausing.contains(v)) { @@ -1274,8 +933,6 @@ class Controller( case WorkerMessage.ReportState(state) => operatorStateMap(workerToOperator(sender)) match { - case PrincipalState.CollectingBreakpoints => - handleWorkerStateReportsInCollBreakpoints(state) case _ => handleWorkerStateReportsInPausing(state) } @@ -1288,8 +945,6 @@ class Controller( disallowActorRefRelatedMessages orElse { case LogErrorToFrontEnd(err: WorkflowRuntimeError) => controllerLogger.logError(err) - case KillAndRecover => - killAndRecoverStage() case QueryStatistics => operatorToWorkerLayers.keys.foreach(opIdentifier => { operatorToWorkerLayers(opIdentifier).foreach(l => { @@ -1344,16 +999,6 @@ class Controller( case scala.util.Failure(t) => throw t } - case PassBreakpointTo(id: String, breakpoint: GlobalBreakpoint) => - val opTag = OperatorIdentifier(tag, id) - operatorToGlobalBreakpoints(opTag)(breakpoint.id) = breakpoint - workflow - .operators(opTag) - .assignBreakpoint( - operatorToWorkerLayers(opTag), - operatorToWorkerStateMap(opTag), - breakpoint - ) case msg => stash() } } diff --git a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/principal/Principal.scala b/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/principal/Principal.scala deleted file mode 100644 index bff16641ac2..00000000000 --- a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/principal/Principal.scala +++ /dev/null @@ -1,796 +0,0 @@ -package edu.uci.ics.amber.engine.architecture.principal - -import akka.actor.{ActorPath, ActorRef, Address, Cancellable, Props} -import akka.pattern.ask -import akka.util.Timeout -import com.google.common.base.Stopwatch -import com.softwaremill.macwire.wire -import edu.uci.ics.amber.clustering.ClusterListener.GetAvailableNodeAddresses -import edu.uci.ics.amber.engine.architecture.breakpoint.FaultedTuple -import edu.uci.ics.amber.engine.architecture.breakpoint.globalbreakpoint.GlobalBreakpoint -import edu.uci.ics.amber.engine.architecture.common.WorkflowActor -import edu.uci.ics.amber.engine.architecture.deploysemantics.layer.WorkerLayer -import edu.uci.ics.amber.engine.architecture.linksemantics.LinkStrategy -import edu.uci.ics.amber.engine.architecture.messaginglayer.NetworkCommunicationActor.RegisterActorRef -import edu.uci.ics.amber.engine.architecture.worker.{WorkerState, WorkerStatistics} -import edu.uci.ics.amber.engine.common.amberexception.WorkflowRuntimeException -import edu.uci.ics.amber.engine.common.ambermessage.ControlMessage._ -import edu.uci.ics.amber.engine.common.ambermessage.ControllerMessage.ReportGlobalBreakpointTriggered -import edu.uci.ics.amber.engine.common.ambermessage.PrincipalMessage.{AssignBreakpoint, _} -import edu.uci.ics.amber.engine.common.ambermessage.StateMessage._ -import edu.uci.ics.amber.engine.common.ambermessage.{PrincipalMessage, WorkerMessage} -import edu.uci.ics.amber.engine.common.ambertag.{AmberTag, LayerTag, WorkerTag} -import edu.uci.ics.amber.engine.common.rpc.AsyncRPCHandlerInitializer -import edu.uci.ics.amber.engine.common.tuple.ITuple -import edu.uci.ics.amber.engine.common.{ - AdvancedMessageSending, - Constants, - TableMetadata, - WorkflowLogger -} -import edu.uci.ics.amber.engine.faulttolerance.recovery.RecoveryPacket -import edu.uci.ics.amber.engine.operators.OpExecConfig -import akka.actor.{ - Actor, - ActorLogging, - ActorPath, - ActorRef, - Address, - Cancellable, - PoisonPill, - Props, - Stash -} -import akka.event.LoggingAdapter -import akka.util.Timeout -import akka.pattern.after -import akka.pattern.ask -import com.google.common.base.Stopwatch -import com.softwaremill.macwire.wire -import edu.uci.ics.amber.engine.architecture.common.WorkflowActor -import edu.uci.ics.amber.engine.architecture.messaginglayer.NetworkCommunicationActor -import com.typesafe.scalalogging.{LazyLogging, Logger} -import edu.uci.ics.amber.engine.architecture.controller.ControllerEvent.ErrorOccurred -import edu.uci.ics.amber.engine.architecture.messaginglayer.NetworkCommunicationActor.RegisterActorRef -import edu.uci.ics.amber.engine.common.ambermessage.WorkerMessage.{ - ReportWorkerPartialCompleted, - ReportedQueriedBreakpoint, - ReportedTriggeredBreakpoints, - Reset -} -import edu.uci.ics.amber.engine.common.ambertag.neo.VirtualIdentity.WorkerActorVirtualIdentity -import edu.uci.ics.amber.engine.common.rpc.AsyncRPCHandlerInitializer -import edu.uci.ics.amber.error.WorkflowRuntimeError - -import scala.collection.mutable -import scala.collection.mutable.ArrayBuffer -import scala.concurrent.duration._ -import scala.concurrent.{Await, ExecutionContext} - -object Principal { - def props(metadata: OpExecConfig, parentNetworkCommunicationActorRef: ActorRef): Props = - Props(new Principal(metadata, parentNetworkCommunicationActorRef)) -} - -class Principal(val metadata: OpExecConfig, parentNetworkCommunicationActorRef: ActorRef) - extends WorkflowActor( - WorkerActorVirtualIdentity(metadata.tag.getGlobalIdentity), - parentNetworkCommunicationActorRef - ) { - implicit val ec: ExecutionContext = context.dispatcher - implicit val timeout: Timeout = 5.seconds - - lazy val rpcHandlerInitializer = wire[AsyncRPCHandlerInitializer] - - private def errorLogAction(err: WorkflowRuntimeError): Unit = { - context.parent ! LogErrorToFrontEnd(err) - } - val errorLogger = WorkflowLogger(s"Principal-${metadata.tag.getGlobalIdentity}-Logger") - errorLogger.setErrorLogAction(errorLogAction) - val tau: FiniteDuration = Constants.defaultTau - var workerLayers: Array[WorkerLayer] = _ - var workerEdges: Array[LinkStrategy] = _ - // var layerDependencies: mutable.HashMap[String, mutable.HashSet[String]] = _ - var workerStateMap: mutable.AnyRefMap[ActorRef, WorkerState.Value] = _ - var workerStatisticsMap: mutable.AnyRefMap[ActorRef, WorkerStatistics] = _ - var workerSinkResultMap = new mutable.AnyRefMap[ActorRef, List[ITuple]] - // var layerMetadata: Array[TableMetadata] = _ - var isUserPaused = false - var globalBreakpoints = new mutable.AnyRefMap[String, GlobalBreakpoint] - var periodicallyAskHandle: Cancellable = _ - var workersTriggeredBreakpoint: Iterable[ActorRef] = _ - var layerCompletedCounter: mutable.HashMap[LayerTag, Int] = _ - val timer = Stopwatch.createUnstarted(); - val stage1Timer = Stopwatch.createUnstarted(); - val stage2Timer = Stopwatch.createUnstarted(); - var receivedRecoveryInformation = new mutable.HashMap[AmberTag, (Long, Long)]() - val receivedTuples = new mutable.ArrayBuffer[(ITuple, ActorPath)]() - - def allWorkerStates: Iterable[WorkerState.Value] = workerStateMap.values - def allWorkers: Iterable[ActorRef] = workerStateMap.keys - def unCompletedWorkerStates: Iterable[WorkerState.Value] = - workerStateMap.filter(x => x._2 != WorkerState.Completed).values - def unCompletedWorkers: Iterable[ActorRef] = - workerStateMap.filter(x => x._2 != WorkerState.Completed).keys - def availableNodes: Array[Address] = - Await - .result(context.actorSelection("/user/cluster-info") ? GetAvailableNodeAddresses, 5.seconds) - .asInstanceOf[Array[Address]] - - private def setWorkerStatistics(worker: ActorRef, workerStatistics: WorkerStatistics): Unit = { - workerStatisticsMap.update(worker, workerStatistics) - } - - // the input count is the sum of the input counts of the first-layer actors -// private def aggregateWorkerInputRowCount(): Long = { -// workerStatisticsMap -// .filter(e => workerLayers.head.layer.contains(e._1)) -// .map(e => e._2.inputRowCount) -// .sum -// } - - // the output count is the sum of the output counts of the last-layer actors -// private def aggregateWorkerOutputRowCount(): Long = { -// workerStatisticsMap -// .filter(e => workerLayers.last.layer.contains(e._1)) -// .map(e => e._2.outputRowCount) -// .sum -// } - - private def setWorkerState(worker: ActorRef, state: WorkerState.Value): Boolean = { - workerStateMap(worker) = state - true - } - - final def whenAllUncompletedWorkersBecome(state: WorkerState.Value): Boolean = - unCompletedWorkerStates.forall(_ == state) - final def whenAllWorkersCompleted: Boolean = allWorkerStates.forall(_ == WorkerState.Completed) - final def safeRemoveAskHandle(): Unit = { - if (periodicallyAskHandle != null) { - periodicallyAskHandle.cancel() - periodicallyAskHandle = null - } - } - - final def resetAll(): Unit = { - workerLayers = null - workerEdges = null - //layerDependencies = null - workerStateMap = null - //layerMetadata = null - isUserPaused = false - safeRemoveAskHandle() - periodicallyAskHandle = null - workersTriggeredBreakpoint = null - layerCompletedCounter = null - globalBreakpoints.foreach(_._2.reset()) - timer.reset() - stage1Timer.reset() - stage2Timer.reset() - context.become(receive) - } - - final def ready: Receive = { - disallowActorRefRelatedMessages orElse [Any, Unit] { - case RecoveryPacket(amberTag, seq1, seq2) => - receivedRecoveryInformation(amberTag) = (seq1, seq2) - case Start => - // sender ! Ack - // allWorkers.foreach(worker => - // AdvancedMessageSending.nonBlockingAskWithRetry(worker, Start, 10, 0) - // ) - case WorkerMessage.ReportState(state) => - // setWorkerState(sender, state) - // state match { - // case WorkerState.Running => - // context.parent ! ReportState(PrincipalState.Running) - // context.become(running) - // timer.start() - // stage1Timer.start() - // unstashAll() - // case WorkerState.Paused => - // if (whenAllUncompletedWorkersBecome(WorkerState.Paused)) { - // safeRemoveAskHandle() - // context.parent ! ReportState(PrincipalState.Paused) - // context.become(paused) - // unstashAll() - // } - // case _ => //throw new AmberException("Invalid worker state received!") - // } - case WorkerMessage.ReportStatistics(statistics) => - // setWorkerStatistics(sender, statistics) - // context.parent ! PrincipalMessage.ReportStatistics( - // PrincipalStatistics( - // PrincipalState.Ready, - // aggregateWorkerInputRowCount(), - // aggregateWorkerOutputRowCount() - // ) - // ) - case StashOutput => - sender ! Ack - allWorkers.foreach(worker => - AdvancedMessageSending.nonBlockingAskWithRetry(worker, StashOutput, 10, 0) - ) - case ReleaseOutput => - sender ! Ack - allWorkers.foreach(worker => - AdvancedMessageSending.nonBlockingAskWithRetry(worker, ReleaseOutput, 10, 0) - ) - case GetInputLayer => sender ! workerLayers.head.clone() - case GetOutputLayer => sender ! workerLayers.last.clone() - case QueryState => sender ! ReportState(PrincipalState.Ready) - case QueryStatistics => - this.allWorkers.foreach(worker => worker ! QueryStatistics) - case Resume => context.parent ! ReportState(PrincipalState.Ready) - case AssignBreakpoint(breakpoint) => - globalBreakpoints(breakpoint.id) = breakpoint - metadata.assignBreakpoint(workerLayers, workerStateMap, breakpoint) - sender ! Ack - case Pause => - allWorkers.foreach(worker => worker ! Pause) - safeRemoveAskHandle() - periodicallyAskHandle = - context.system.scheduler.schedule(0.milliseconds, 30.seconds, self, EnforceStateCheck) - context.become(pausing) - unstashAll() - case msg => stash() - } - } - - final def running: Receive = { - case RecoveryPacket(amberTag, seq1, seq2) => - receivedRecoveryInformation(amberTag) = (seq1, seq2) - case WorkerMessage.ReportState(state) => -// log.info("running: " + sender + " to " + state) -// if (setWorkerState(sender, state)) { -// state match { -// case WorkerState.LocalBreakpointTriggered => -// if (whenAllUncompletedWorkersBecome(WorkerState.LocalBreakpointTriggered)) { -// //only one worker and it triggered breakpoint -// safeRemoveAskHandle() -// periodicallyAskHandle = context.system.scheduler.schedule( -// 0.milliseconds, -// 30.seconds, -// self, -// EnforceStateCheck -// ) -// workersTriggeredBreakpoint = allWorkers -// context.parent ! ReportState(PrincipalState.CollectingBreakpoints) -// context.become(collectingBreakpoints) -// } else { -// //no tau involved since we know a very small tau works best -// if (!stage2Timer.isRunning) { -// stage2Timer.start() -// } -// if (stage1Timer.isRunning) { -// stage1Timer.stop() -// } -// context.system.scheduler -// .scheduleOnce(tau, () => unCompletedWorkers.foreach(worker => worker ! Pause)) -// safeRemoveAskHandle() -// periodicallyAskHandle = -// context.system.scheduler.schedule(30.seconds, 30.seconds, self, EnforceStateCheck) -// context.become(pausing) -// unstashAll() -// } -// case WorkerState.Paused => -// if (whenAllWorkersCompleted) { -// safeRemoveAskHandle() -// context.parent ! ReportState(PrincipalState.Completed) -// context.become(completed) -// unstashAll() -// } else if (whenAllUncompletedWorkersBecome(WorkerState.Paused)) { -// safeRemoveAskHandle() -// context.parent ! ReportState(PrincipalState.Paused) -// context.become(paused) -// unstashAll() -// } else if ( -// unCompletedWorkerStates -// .forall(x => x == WorkerState.Paused || x == WorkerState.LocalBreakpointTriggered) -// ) { -// workersTriggeredBreakpoint = -// workerStateMap.filter(_._2 == WorkerState.LocalBreakpointTriggered).keys -// safeRemoveAskHandle() -// periodicallyAskHandle = context.system.scheduler.schedule( -// 1.milliseconds, -// 30.seconds, -// self, -// EnforceStateCheck -// ) -// context.parent ! ReportState(PrincipalState.CollectingBreakpoints) -// context.become(collectingBreakpoints) -// unstashAll() -// } -// case WorkerState.Completed => -// if (whenAllWorkersCompleted) { -// if (timer.isRunning) { -// timer.stop() -// } -// log.info(metadata.tag.toString + " completed! Time Elapsed: " + timer.toString()) -// context.parent ! ReportState(PrincipalState.Completed) -// context.become(completed) -// unstashAll() -// } -// case _ => //skip others for now -// } -// } - case WorkerMessage.ReportStatistics(statistics) => -// setWorkerStatistics(sender, statistics) -// context.parent ! PrincipalMessage.ReportStatistics( -// PrincipalStatistics( -// PrincipalState.Running, -// aggregateWorkerInputRowCount(), -// aggregateWorkerOutputRowCount() -// ) -// ) - case Pause => - //single point pause: pause itself -// if (sender != self) { -// isUserPaused = true -// } -// allWorkers.foreach(worker => worker ! Pause) -// safeRemoveAskHandle() -// periodicallyAskHandle = -// context.system.scheduler.schedule(30.seconds, 30.seconds, self, EnforceStateCheck) -// context.become(pausing) -// unstashAll() - case ReportWorkerPartialCompleted(worker, layer) => -// sender ! Ack -// AdvancedMessageSending.nonBlockingAskWithRetry( -// context.parent, -// ReportPrincipalPartialCompleted(worker, layer), -// 10, -// 0 -// ) -// if (layerCompletedCounter.contains(layer)) { -// layerCompletedCounter(layer) -= 1 -// if (layerCompletedCounter(layer) == 0) { -// layerCompletedCounter -= layer -// AdvancedMessageSending.nonBlockingAskWithRetry( -// context.parent, -// ReportPrincipalPartialCompleted(metadata.tag, layer), -// 10, -// 0 -// ) -// } -// } - case StashOutput => - sender ! Ack - allWorkers.foreach(worker => - AdvancedMessageSending.nonBlockingAskWithRetry(worker, StashOutput, 10, 0) - ) - case ReleaseOutput => -// sender ! Ack -// allWorkers.foreach(worker => -// AdvancedMessageSending.nonBlockingAskWithRetry(worker, ReleaseOutput, 10, 0) -// ) - case GetInputLayer => sender ! workerLayers.head.clone() - case GetOutputLayer => sender ! workerLayers.last.clone() - case Resume => context.parent ! ReportState(PrincipalState.Running) - case QueryState => sender ! ReportState(PrincipalState.Running) - case QueryStatistics => - this.allWorkers.foreach(worker => worker ! QueryStatistics) - case msg => - //log.info("stashing: "+ msg) - stash() - } - -// final lazy val allowedStatesOnPausing: Set[WorkerState.Value] = -// Set(WorkerState.Completed, WorkerState.Paused, WorkerState.LocalBreakpointTriggered) - - final def pausing: Receive = { - disallowActorRefRelatedMessages orElse [Any, Unit] { - case RecoveryPacket(amberTag, seq1, seq2) => - receivedRecoveryInformation(amberTag) = (seq1, seq2) - case EnforceStateCheck => - // for ((k, v) <- workerStateMap) { - // if (!allowedStatesOnPausing.contains(v)) { - // k ! QueryState - // } - // } - case reportCurrentTuple: WorkerMessage.ReportCurrentProcessingTuple => - // receivedTuples.append((reportCurrentTuple.tuple, reportCurrentTuple.workerID)) - case WorkerMessage.ReportState(state) => - // log.info("pausing: " + sender + " to " + state) - // if (setWorkerState(sender, state)) { - // if (whenAllWorkersCompleted) { - // safeRemoveAskHandle() - // context.parent ! ReportState(PrincipalState.Completed) - // context.become(completed) - // unstashAll() - // } else if (whenAllUncompletedWorkersBecome(WorkerState.Paused)) { - // safeRemoveAskHandle() - // context.parent ! ReportCurrentProcessingTuple( - // this.metadata.tag.operator, - // receivedTuples.toArray - // ) - // receivedTuples.clear() - // context.parent ! ReportState(PrincipalState.Paused) - // context.become(paused) - // unstashAll() - // } else if ( - // unCompletedWorkerStates - // .forall(x => x == WorkerState.Paused || x == WorkerState.LocalBreakpointTriggered) - // ) { - // workersTriggeredBreakpoint = - // workerStateMap.filter(_._2 == WorkerState.LocalBreakpointTriggered).keys - // safeRemoveAskHandle() - // periodicallyAskHandle = - // context.system.scheduler.schedule(1.milliseconds, 30.seconds, self, EnforceStateCheck) - // context.parent ! ReportState(PrincipalState.CollectingBreakpoints) - // context.become(collectingBreakpoints) - // unstashAll() - // } - // } - case WorkerMessage.ReportStatistics(statistics) => - // setWorkerStatistics(sender, statistics) - // context.parent ! PrincipalMessage.ReportStatistics( - // PrincipalStatistics( - // PrincipalState.Pausing, - // aggregateWorkerInputRowCount(), - // aggregateWorkerOutputRowCount() - // ) - // ) - case QueryState => sender ! ReportState(PrincipalState.Pausing) - case QueryStatistics => - this.allWorkers.foreach(worker => worker ! QueryStatistics) - case Pause => - if (sender != self) { - isUserPaused = true - } - case msg => - //log.info("stashing: "+ msg) - stash() - } - } - - final def collectingBreakpoints: Receive = { - disallowActorRefRelatedMessages orElse [Any, Unit] { - case RecoveryPacket(amberTag, seq1, seq2) => - receivedRecoveryInformation(amberTag) = (seq1, seq2) - case EnforceStateCheck => - // workersTriggeredBreakpoint.foreach(x => x ! QueryTriggeredBreakpoints) //query all - case WorkerMessage.ReportState(state) => - //log.info("collecting: "+ sender +" to "+ state) - // if (setWorkerState(sender, state)) { - // if (unCompletedWorkerStates.forall(_ == WorkerState.Paused)) { - ////all breakpoint resolved, it's safe to report to controller and then Pause(on triggered, or user paused) else Resume - // val map = new mutable.HashMap[(ActorRef, FaultedTuple), ArrayBuffer[String]] - // for (i <- globalBreakpoints.values.filter(_.isTriggered)) { - // isUserPaused = true //upgrade pause - // i.report(map) - // } - // safeRemoveAskHandle() - // context.become(paused) - // unstashAll() - // if (!isUserPaused) { - // log.info("no global breakpoint triggered, continue") - // self ! Resume - // } else { - // context.parent ! ReportGlobalBreakpointTriggered(map, this.metadata.tag.operator) - // context.parent ! ReportState(PrincipalState.Paused) - // log.info( - // "user paused or global breakpoint triggered, pause. Stage1 cost = " + stage1Timer - // .toString() + " Stage2 cost =" + stage2Timer.toString() - // ) - // } - // if (stage2Timer.isRunning) { - // stage2Timer.stop() - // } - // if (!stage1Timer.isRunning) { - // stage1Timer.start() - // } - // } - // } - case WorkerMessage.ReportStatistics(statistics) => - // setWorkerStatistics(sender, statistics) - // context.parent ! PrincipalMessage.ReportStatistics( - // PrincipalStatistics( - // PrincipalState.CollectingBreakpoints, - // aggregateWorkerInputRowCount(), - // aggregateWorkerOutputRowCount() - // ) - // ) - case ReportedTriggeredBreakpoints(bps) => - // bps.foreach(x => { - // val bp = globalBreakpoints(x.id) - // bp.accept(sender, x) - // if (bp.needCollecting) { - ////is not fully collected - // bp.collect() - // } else if (bp.isRepartitionRequired) { - ////fully collected, but need repartition (e.g. count not reach target number) - ////OR need Reset - // metadata.assignBreakpoint(workerLayers, workerStateMap, bp) - // } else if (bp.isCompleted) { - ////fully collected and reach the target - // bp.remove() - // } - // }) - case ReportedQueriedBreakpoint(bp) => - // val gbp = globalBreakpoints(bp.id) - // if (gbp.accept(sender, bp) && !gbp.needCollecting) { - // if (gbp.isRepartitionRequired) { - ////fully collected, but need repartition (count not reach target number) - // metadata.assignBreakpoint(workerLayers, workerStateMap, gbp) - // } else if (gbp.isCompleted) { - ////fully collected and reach the target - // gbp.remove() - // } - // } - case GetInputLayer => sender ! workerLayers.head.clone() - case GetOutputLayer => sender ! workerLayers.last.clone() - case Pause => - if (sender != self) { - isUserPaused = true - } - case msg => - //log.info("stashing: "+ msg) - stash() - } - } - -// final lazy val allowedStatesOnResuming: Set[WorkerState.Value] = -// Set(WorkerState.Running, WorkerState.Ready, WorkerState.Completed) - - final def resuming: Receive = { - disallowActorRefRelatedMessages orElse [Any, Unit] { - case RecoveryPacket(amberTag, seq1, seq2) => - receivedRecoveryInformation(amberTag) = (seq1, seq2) - case EnforceStateCheck => - //for ((k, v) <- workerStateMap) { - // if (!allowedStatesOnResuming.contains(v)) { - // k ! QueryState - // } - //} - case WorkerMessage.ReportState(state) => - //log.info("resuming: "+ sender +" to "+ state) - //if (!allowedStatesOnResuming.contains(state)) { - // sender ! Resume - //} else if (setWorkerState(sender, state)) { - // if (whenAllWorkersCompleted) { - // safeRemoveAskHandle() - // context.parent ! ReportState(PrincipalState.Completed) - // context.become(completed) - // unstashAll() - // } else if (allWorkerStates.forall(_ != WorkerState.Paused)) { - // safeRemoveAskHandle() - // if (allWorkerStates.exists(_ != WorkerState.Ready)) { - // context.parent ! ReportState(PrincipalState.Running) - // context.become(running) - // } else { - // context.parent ! ReportState(PrincipalState.Ready) - // context.become(ready) - // } - // unstashAll() - // } - //} - case GetInputLayer => sender ! workerLayers.head.clone() - case GetOutputLayer => sender ! workerLayers.last.clone() - case WorkerMessage.ReportStatistics(statistics) => - //setWorkerStatistics(sender, statistics) - //context.parent ! PrincipalMessage.ReportStatistics( - // PrincipalStatistics( - // PrincipalState.Resuming, - // aggregateWorkerInputRowCount(), - // aggregateWorkerOutputRowCount() - // ) - //) - case QueryState => sender ! ReportState(PrincipalState.Resuming) - case QueryStatistics => - this.allWorkers.foreach(worker => worker ! QueryStatistics) - case msg => - //log.info("stashing: "+ msg) - stash() - } - } - - final def paused: Receive = { - disallowActorRefRelatedMessages orElse [Any, Unit] { - case KillAndRecover => -// workerLayers.foreach { x => -// x.layer(0) ! Reset(x.getFirstMetadata, Seq(receivedRecoveryInformation(x.tagForFirst))) -// workerStateMap(x.layer(0)) = WorkerState.Ready -// } -// sender ! Ack -// context.become(pausing) - case RecoveryPacket(amberTag, seq1, seq2) => - receivedRecoveryInformation(amberTag) = receivedRecoveryInformation(amberTag) - case Resume => - // isUserPaused = false //reset - // assert(unCompletedWorkerStates.nonEmpty) - // unCompletedWorkers.foreach(worker => worker ! Resume) - // safeRemoveAskHandle() - // periodicallyAskHandle = - // context.system.scheduler.schedule(30.seconds, 30.seconds, self, EnforceStateCheck) - // context.become(resuming) - // unstashAll() - case AssignBreakpoint(breakpoint) => - // sender ! Ack - // globalBreakpoints(breakpoint.id) = breakpoint - // metadata.assignBreakpoint(workerLayers, workerStateMap, breakpoint) - case GetInputLayer => sender ! workerLayers.head.clone() - case GetOutputLayer => sender ! workerLayers.last.clone() - case Pause => context.parent ! ReportState(PrincipalState.Paused) - case QueryState => sender ! ReportState(PrincipalState.Paused) - case ModifyLogic(newMetadata) => - // sender ! Ack - // log.info("modify logic received by principal, sending to worker") - // this.allWorkers.foreach(worker => worker ! ModifyLogic(newMetadata)) - //// allWorkers.foreach(worker => AdvancedMessageSending.blockingAskWithRetry(worker, ModifyLogic(newMetadata), 3)) - // log.info("modify logic received by principal, sent to worker") - case QueryStatistics => - this.allWorkers.foreach(worker => worker ! QueryStatistics) - case WorkerMessage.ReportStatistics(statistics) => - // setWorkerStatistics(sender, statistics) - // context.parent ! PrincipalMessage.ReportStatistics( - // PrincipalStatistics( - // PrincipalState.Paused, - // aggregateWorkerInputRowCount(), - // aggregateWorkerOutputRowCount() - // ) - // ) - case msg => - //log.info("stashing: "+ msg) - stash() - } - } - - final def completed: Receive = { - disallowActorRefRelatedMessages orElse [Any, Unit] { - case KillAndRecover => -// workerLayers.foreach { x => -// if (receivedRecoveryInformation.contains(x.tagForFirst)) { -// x.layer(0) ! Reset(x.getFirstMetadata, Seq(receivedRecoveryInformation(x.tagForFirst))) -// } else { -// x.layer(0) ! Reset(x.getFirstMetadata, Seq()) -// } -// workerStateMap(x.layer(0)) = WorkerState.Ready -// } -// sender ! Ack -// context.become(pausing) - case RecoveryPacket(amberTag, seq1, seq2) => - receivedRecoveryInformation(amberTag) = (seq1, seq2) - case QueryStatistics => - this.allWorkers.foreach(worker => worker ! QueryStatistics) - case StashOutput => - sender ! Ack - allWorkers.foreach(worker => - AdvancedMessageSending.nonBlockingAskWithRetry(worker, StashOutput, 10, 0) - ) - case ReleaseOutput => - sender ! Ack - allWorkers.foreach(worker => - AdvancedMessageSending.nonBlockingAskWithRetry(worker, ReleaseOutput, 10, 0) - ) - case WorkerMessage.ReportStatistics(statistics) => - //setWorkerStatistics(sender, statistics) - //context.parent ! PrincipalMessage.ReportStatistics( - // PrincipalStatistics( - // PrincipalState.Completed, - // aggregateWorkerInputRowCount(), - // aggregateWorkerOutputRowCount() - // ) - //) - case CollectSinkResults => - //allWorkers.foreach(worker => worker ! CollectSinkResults) - case WorkerMessage.ReportOutputResult(sinkResult) => - //workerSinkResultMap(sender) = sinkResult - //if (workerSinkResultMap.size == allWorkers.size) { - // val collectedResults = mutable.MutableList[ITuple]() - // this.workerSinkResultMap.values.foreach(v => collectedResults ++= v) - // context.parent ! PrincipalMessage.ReportOutputResult(collectedResults.toList) - //} - case GetInputLayer => sender ! workerLayers.head.clone() - case GetOutputLayer => sender ! workerLayers.last.clone() - case msg => - //log.info("received {} from {} after complete",msg,sender) - if (sender == context.parent) { - sender ! ReportState(PrincipalState.Completed) - } - } - } - - final override def receive: Receive = { - disallowActorRefRelatedMessages orElse [Any, Unit] { - case AckedPrincipalInitialization(prev: Array[(OpExecConfig, WorkerLayer)]) => - //workerLayers = metadata.topology.layers - //workerEdges = metadata.topology.links - //val all = availableNodes - //if (workerEdges.isEmpty) { - // workerLayers.foreach(x => x.build(prev, all)) - //} else { - // val inLinks: Map[WorkerLayer, Set[WorkerLayer]] = - // workerEdges.groupBy(x => x.to).map(x => (x._1, x._2.map(_.from).toSet)) - // var currentLayer: Iterable[WorkerLayer] = - // workerEdges.filter(x => workerEdges.forall(_.to != x.from)).map(_.from) - // currentLayer.foreach(x => x.build(prev, all)) - // currentLayer = inLinks.filter(x => x._2.forall(_.isBuilt)).keys - // while (currentLayer.nonEmpty) { - // currentLayer.foreach(x => x.build(inLinks(x).map(y => (null, y)).toArray, all)) - // currentLayer = inLinks.filter(x => !x._1.isBuilt && x._2.forall(_.isBuilt)).keys - // } - //} - //layerCompletedCounter = - // mutable.HashMap(prev.map(x => x._2.tag -> workerLayers.head.layer.length).toSeq: _*) - //workerStateMap = mutable.AnyRefMap( - // workerLayers.flatMap(x => x.layer).map((_, WorkerState.Uninitialized)).toMap.toSeq: _* - //) - //workerStatisticsMap = mutable.AnyRefMap( - // workerLayers - // .flatMap(x => x.layer) - // .map((_, WorkerStatistics(WorkerState.Uninitialized, 0, 0))) - // .toMap - // .toSeq: _* - //) - //workerLayers.foreach { x => - // var i = 0 - // x.layer.foreach { worker => - // val workerTag = WorkerTag(x.tag, i) - // worker ! AckedWorkerInitialization() - // i += 1 - // } - //} - //safeRemoveAskHandle() - //periodicallyAskHandle = - // context.system.scheduler.schedule(30.seconds, 30.seconds, self, EnforceStateCheck) - //context.become(initializing) - //unstashAll() - //sender ! AckWithInformation(metadata) - case QueryState => sender ! ReportState(PrincipalState.Uninitialized) - case QueryStatistics => - //this.allWorkers.foreach(worker => worker ! QueryStatistics) - //sender() ! ReportStatistics( - // PrincipalStatistics( - // PrincipalState.Uninitialized, - // aggregateWorkerInputRowCount(), - // aggregateWorkerOutputRowCount() - // ) - //) - case msg => - //log.info("stashing: "+ msg) - stash() - } - } - - final def initializing: Receive = { - case EnforceStateCheck => -// for ((k, v) <- workerStateMap) { -// if (v != WorkerState.Ready) { -// k ! QueryState -// } -// } - case WorkerMessage.ReportState(state) => -// if (state != WorkerState.Ready) { -// sender ! AckedWorkerInitialization() -// } else if (setWorkerState(sender, state)) { -// if (whenAllUncompletedWorkersBecome(WorkerState.Ready)) { -// safeRemoveAskHandle() -// workerEdges.foreach(x => x.link()) -// globalBreakpoints.values.foreach( -// metadata.assignBreakpoint(workerLayers, workerStateMap, _) -// ) -// allWorkers.foreach(_ ! CheckRecovery) -// context.parent ! ReportState(PrincipalState.Ready) -// context.become(ready) -// unstashAll() -// } -// } - case WorkerMessage.ReportStatistics(statistics) => -// setWorkerStatistics(sender, statistics) -// case QueryState => sender ! ReportState(PrincipalState.Initializing) -// case QueryStatistics => -// this.allWorkers.foreach(worker => worker ! QueryStatistics) -// sender() ! ReportStatistics( -// PrincipalStatistics( -// PrincipalState.Initializing, -// aggregateWorkerInputRowCount(), -// aggregateWorkerOutputRowCount() -// ) -// ) - case msg => - //log.info("stashing: "+ msg) - stash() - } - -} diff --git a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/worker/WorkflowWorker.scala b/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/worker/WorkflowWorker.scala index 60aad2a0ef6..824fca4c0d6 100644 --- a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/worker/WorkflowWorker.scala +++ b/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/worker/WorkflowWorker.scala @@ -151,14 +151,6 @@ class WorkflowWorker( logger.logInfo(s"received register input for ${this.identifier}") sender ! Ack tupleProducer.registerInput(identifier, inputNum) - case LocalBreakpointTriggered() => - workerStateManager.confirmState(Running) - dataProcessor.breakpoints.foreach { brk => - if (brk.isTriggered) - dataProcessor.unhandledFaultedTuples(brk.triggeredTupleId) = - new FaultedTuple(brk.triggeredTuple, brk.triggeredTupleId, brk.isInput) - } - context.parent ! ReportState(WorkerState.LocalBreakpointTriggered) case other => } @@ -225,69 +217,8 @@ class WorkflowWorker( sender ! ReportStatistics( WorkerStatistics(getOldWorkerState, getInputRowCount(), getOutputRowCount()) ) - case QueryTriggeredBreakpoints => - val toReport = dataProcessor.breakpoints.filter(_.isTriggered) - if (toReport.nonEmpty) { - toReport.foreach(_.isReported = true) - sender ! ReportedTriggeredBreakpoints(toReport) - } else { - throw new WorkflowRuntimeException( - WorkflowRuntimeError( - "no triggered local breakpoints but worker in triggered breakpoint state", - "WorkerBase:allowQueryTriggeredBreakpoints", - Map() - ) - ) - } - case QueryBreakpoint(id) => - val toReport = dataProcessor.breakpoints.find(_.id == id) - if (toReport.isDefined) { - toReport.get.isReported = true - context.parent ! ReportedQueriedBreakpoint(toReport.get) - context.parent ! ReportState(WorkerState.LocalBreakpointTriggered) - } case CollectSinkResults => sender ! WorkerMessage.ReportOutputResult(this.getResultTuples().toList) - case AssignBreakpoint(bp) => - sender ! Ack - dataProcessor.registerBreakpoint(bp) -// if (!dataProcessor.breakpoints.exists(_.isDirty)) { -// onPaused() //back to paused -// context.unbecome() -// unstashAll() -// } - case RemoveBreakpoint(id) => - sender ! Ack - dataProcessor.removeBreakpoint(id) -// if (!dataProcessor.breakpoints.exists(_.isDirty)) { -// onPaused() //back to paused -// context.unbecome() -// unstashAll() -// } - case SkipTuple(f) => - workerStateManager.confirmState(Paused) - sender ! Ack - if (!receivedFaultedTupleIds.contains(f.id)) { - receivedFaultedTupleIds.add(f.id) - dataProcessor.unhandledFaultedTuples.remove(f.id) - onSkipTuple(f) - } - case ModifyTuple(f) => - workerStateManager.confirmState(Paused) - sender ! Ack - if (!receivedFaultedTupleIds.contains(f.id)) { - receivedFaultedTupleIds.add(f.id) - dataProcessor.unhandledFaultedTuples.remove(f.id) - onModifyTuple(f) - } - case ResumeTuple(f) => - workerStateManager.confirmState(Paused) - sender ! Ack - if (!receivedFaultedTupleIds.contains(f.id)) { - receivedFaultedTupleIds.add(f.id) - dataProcessor.unhandledFaultedTuples.remove(f.id) - onResumeTuple(f) - } case AddDataSendingPolicy(policy) => sender ! Ack // send message to receivers to add this worker to their expected inputs diff --git a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/worker/neo/DataProcessor.scala b/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/worker/neo/DataProcessor.scala index e22cd06d428..25b23446735 100644 --- a/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/worker/neo/DataProcessor.scala +++ b/core/amber/src/main/scala/edu/uci/ics/amber/engine/architecture/worker/neo/DataProcessor.scala @@ -3,7 +3,6 @@ package edu.uci.ics.amber.engine.architecture.worker.neo import java.util.concurrent.Executors import com.typesafe.scalalogging.LazyLogging -import edu.uci.ics.amber.engine.architecture.breakpoint.localbreakpoint.ExceptionBreakpoint import edu.uci.ics.amber.engine.architecture.messaginglayer.{ ControlOutputPort, TupleToBatchConverter @@ -22,8 +21,7 @@ class DataProcessor( // dependencies: controlOutputChannel: ControlOutputPort, // to send controls to main thread batchProducer: TupleToBatchConverter, // to send output tuples pauseManager: PauseManager // to pause/resume -) extends BreakpointSupport - with WorkerInternalQueue { // TODO: make breakpointSupport as a module +) extends WorkerInternalQueue { // TODO: make breakpointSupport as a module protected val logger: WorkflowLogger = WorkflowLogger("DataProcessor") // dp thread stats: @@ -135,22 +133,8 @@ class DataProcessor( // dependencies: controlOutputChannel.sendTo(VirtualIdentity.Self, ExecutionCompleted()) } - // For compatibility, we use old breakpoint handling logic - // TODO: remove this when we refactor breakpoints - private[this] def assignExceptionBreakpoint( - faultedTuple: ITuple, - e: Exception, - isInput: Boolean - ): Unit = { - breakpoints(0).triggeredTuple = faultedTuple - breakpoints(0).asInstanceOf[ExceptionBreakpoint].error = e - breakpoints(0).triggeredTupleId = outputTupleCount - breakpoints(0).isInput = isInput - } - private[this] def handleOperatorException(e: Exception, isInput: Boolean): Unit = { pauseManager.pause() - assignExceptionBreakpoint(currentInputTuple.left.getOrElse(null), e, isInput) controlOutputChannel.sendTo(VirtualIdentity.Self, LocalBreakpointTriggered()) } diff --git a/core/amber/src/main/scala/edu/uci/ics/amber/engine/operators/OpExecConfig.scala b/core/amber/src/main/scala/edu/uci/ics/amber/engine/operators/OpExecConfig.scala index dcbb4e37843..9ce54a064a7 100644 --- a/core/amber/src/main/scala/edu/uci/ics/amber/engine/operators/OpExecConfig.scala +++ b/core/amber/src/main/scala/edu/uci/ics/amber/engine/operators/OpExecConfig.scala @@ -10,6 +10,7 @@ import edu.uci.ics.amber.engine.common.tuple.ITuple import akka.actor.ActorRef import akka.event.LoggingAdapter import akka.util.Timeout +import edu.uci.ics.amber.engine.common.ambertag.neo.VirtualIdentity.ActorVirtualIdentity import scala.collection.mutable import scala.concurrent.ExecutionContext @@ -50,10 +51,6 @@ abstract class OpExecConfig(val tag: OperatorIdentifier) extends Serializable { def getShuffleHashFunction(layerTag: LayerTag): ITuple => Int = ??? - def assignBreakpoint( - topology: Array[WorkerLayer], - states: mutable.AnyRefMap[ActorRef, WorkerState.Value], - breakpoint: GlobalBreakpoint - )(implicit timeout: Timeout, ec: ExecutionContext) + def assignBreakpoint(breakpoint: GlobalBreakpoint): Array[ActorVirtualIdentity] } diff --git a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/operators/OneToOneOpExecConfig.scala b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/operators/OneToOneOpExecConfig.scala index 1c8989f3ec2..930619289f3 100644 --- a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/operators/OneToOneOpExecConfig.scala +++ b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/operators/OneToOneOpExecConfig.scala @@ -9,6 +9,7 @@ import edu.uci.ics.amber.engine.architecture.deploysemantics.deploystrategy.Roun import edu.uci.ics.amber.engine.architecture.deploysemantics.layer.WorkerLayer import edu.uci.ics.amber.engine.architecture.worker.WorkerState import edu.uci.ics.amber.engine.common.Constants +import edu.uci.ics.amber.engine.common.ambertag.neo.VirtualIdentity.ActorVirtualIdentity import edu.uci.ics.amber.engine.common.ambertag.{LayerTag, OperatorIdentifier} import edu.uci.ics.amber.engine.operators.OpExecConfig @@ -37,12 +38,9 @@ class OneToOneOpExecConfig( } override def assignBreakpoint( - topology: Array[WorkerLayer], - states: mutable.AnyRefMap[ActorRef, WorkerState.Value], breakpoint: GlobalBreakpoint - )(implicit timeout: Timeout, ec: ExecutionContext): Unit = { - breakpoint.partition( - topology(0).layer.filter(states(_) != WorkerState.Completed) - ) + ): Array[ActorVirtualIdentity] = { + // TODO: take worker states into account + topology.layers(0).identifiers } } diff --git a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/operators/aggregate/AggregateOpExecConfig.scala b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/operators/aggregate/AggregateOpExecConfig.scala index 0ea24098123..d572af4320e 100644 --- a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/operators/aggregate/AggregateOpExecConfig.scala +++ b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/operators/aggregate/AggregateOpExecConfig.scala @@ -17,6 +17,7 @@ import edu.uci.ics.amber.engine.architecture.deploysemantics.layer.WorkerLayer import edu.uci.ics.amber.engine.architecture.linksemantics.{AllToOne, HashBasedShuffle} import edu.uci.ics.amber.engine.architecture.worker.WorkerState import edu.uci.ics.amber.engine.common.Constants +import edu.uci.ics.amber.engine.common.ambertag.neo.VirtualIdentity.ActorVirtualIdentity import edu.uci.ics.amber.engine.common.ambertag.{LayerTag, OperatorIdentifier} import edu.uci.ics.amber.engine.operators.OpExecConfig import edu.uci.ics.texera.workflow.common.tuple.Tuple @@ -92,12 +93,11 @@ class AggregateOpExecConfig[P <: AnyRef]( ) } } + override def assignBreakpoint( - topology: Array[WorkerLayer], - states: mutable.AnyRefMap[ActorRef, WorkerState.Value], breakpoint: GlobalBreakpoint - )(implicit timeout: Timeout, ec: ExecutionContext): Unit = { - breakpoint.partition(topology(0).layer.filter(states(_) != WorkerState.Completed)) + ): Array[ActorVirtualIdentity] = { + topology.layers(0).identifiers } } diff --git a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/operators/mlmodel/MLModelOpExecConfig.scala b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/operators/mlmodel/MLModelOpExecConfig.scala index 1ee238fdec4..ec84708e460 100644 --- a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/operators/mlmodel/MLModelOpExecConfig.scala +++ b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/operators/mlmodel/MLModelOpExecConfig.scala @@ -9,6 +9,7 @@ import edu.uci.ics.amber.engine.architecture.deploysemantics.deploystrategy.Roun import edu.uci.ics.amber.engine.architecture.deploysemantics.layer.WorkerLayer import edu.uci.ics.amber.engine.architecture.worker.WorkerState import edu.uci.ics.amber.engine.common.Constants +import edu.uci.ics.amber.engine.common.ambertag.neo.VirtualIdentity.ActorVirtualIdentity import edu.uci.ics.amber.engine.common.ambertag.{LayerTag, OperatorIdentifier} import edu.uci.ics.amber.engine.operators.OpExecConfig import edu.uci.ics.texera.workflow.common.operators.OperatorExecutor @@ -37,14 +38,10 @@ class MLModelOpExecConfig( Map() ) } + override def assignBreakpoint( - topology: Array[WorkerLayer], - states: mutable.AnyRefMap[ActorRef, WorkerState.Value], breakpoint: GlobalBreakpoint - )(implicit timeout: Timeout, ec: ExecutionContext): Unit = { - breakpoint.partition( - topology(0).layer.filter(states(_) != WorkerState.Completed) - ) + ): Array[ActorVirtualIdentity] = { + topology.layers(0).identifiers } - } diff --git a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/workflow/WorkflowCompiler.scala b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/workflow/WorkflowCompiler.scala index 93d11f4dc50..395fdff9e5f 100644 --- a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/workflow/WorkflowCompiler.scala +++ b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/common/workflow/WorkflowCompiler.scala @@ -1,10 +1,6 @@ package edu.uci.ics.texera.workflow.common.workflow import akka.actor.ActorRef -import edu.uci.ics.amber.engine.architecture.breakpoint.globalbreakpoint.{ - ConditionalGlobalBreakpoint, - CountGlobalBreakpoint -} import edu.uci.ics.amber.engine.architecture.controller.Workflow import edu.uci.ics.amber.engine.common.ambermessage.ControllerMessage.PassBreakpointTo import edu.uci.ics.amber.engine.common.ambertag.OperatorIdentifier @@ -95,21 +91,21 @@ class WorkflowCompiler(val workflowInfo: WorkflowInfo, val context: WorkflowCont case BreakpointCondition.NOT_CONTAINS => tuple => !tuple.getField(column).toString.trim.contains(conditionBp.value) } - controller ! PassBreakpointTo( - operatorID, - new ConditionalGlobalBreakpoint( - breakpointID, - tuple => { - val texeraTuple = tuple.asInstanceOf[Tuple] - predicate.apply(texeraTuple) - } - ) - ) +// controller ! PassBreakpointTo( +// operatorID, +// new ConditionalGlobalBreakpoint( +// breakpointID, +// tuple => { +// val texeraTuple = tuple.asInstanceOf[Tuple] +// predicate.apply(texeraTuple) +// } +// ) +// ) case countBp: CountBreakpoint => - controller ! PassBreakpointTo( - operatorID, - new CountGlobalBreakpoint("breakpointID", countBp.count) - ) +// controller ! PassBreakpointTo( +// operatorID, +// new CountGlobalBreakpoint("breakpointID", countBp.count) +// ) } } diff --git a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/localscan/LocalCsvFileScanOpExecConfig.scala b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/localscan/LocalCsvFileScanOpExecConfig.scala index 8ad534493da..a2e9826335d 100644 --- a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/localscan/LocalCsvFileScanOpExecConfig.scala +++ b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/localscan/LocalCsvFileScanOpExecConfig.scala @@ -10,6 +10,7 @@ import edu.uci.ics.amber.engine.architecture.deploysemantics.deploymentfilter.Us import edu.uci.ics.amber.engine.architecture.deploysemantics.deploystrategy.RoundRobinDeployment import edu.uci.ics.amber.engine.architecture.deploysemantics.layer.WorkerLayer import edu.uci.ics.amber.engine.architecture.worker.WorkerState +import edu.uci.ics.amber.engine.common.ambertag.neo.VirtualIdentity.ActorVirtualIdentity import edu.uci.ics.amber.engine.common.ambertag.{LayerTag, OperatorIdentifier} import edu.uci.ics.amber.engine.operators.OpExecConfig import edu.uci.ics.texera.workflow.common.tuple.schema.Schema @@ -53,13 +54,10 @@ class LocalCsvFileScanOpExecConfig( Map() ) } - override def assignBreakpoint( - topology: Array[WorkerLayer], - states: mutable.AnyRefMap[ActorRef, WorkerState.Value], breakpoint: GlobalBreakpoint - )(implicit timeout: Timeout, ec: ExecutionContext): Unit = { - breakpoint.partition(topology(0).layer.filter(states(_) != WorkerState.Completed)) + ): Array[ActorVirtualIdentity] = { + topology.layers(0).identifiers } } diff --git a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/pythonUDF/PythonUDFOpExecConfig.scala b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/pythonUDF/PythonUDFOpExecConfig.scala index 4e9f92afd21..8816cea9662 100644 --- a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/pythonUDF/PythonUDFOpExecConfig.scala +++ b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/pythonUDF/PythonUDFOpExecConfig.scala @@ -10,6 +10,7 @@ import edu.uci.ics.amber.engine.architecture.deploysemantics.deploymentfilter.Fo import edu.uci.ics.amber.engine.architecture.deploysemantics.deploystrategy.RoundRobinDeployment import edu.uci.ics.amber.engine.architecture.deploysemantics.layer.WorkerLayer import edu.uci.ics.amber.engine.architecture.worker.WorkerState +import edu.uci.ics.amber.engine.common.ambertag.neo.VirtualIdentity.ActorVirtualIdentity import edu.uci.ics.amber.engine.common.ambertag.{LayerTag, OperatorIdentifier} import edu.uci.ics.amber.engine.operators.OpExecConfig import edu.uci.ics.texera.workflow.common.tuple.schema.Attribute @@ -52,11 +53,9 @@ class PythonUDFOpExecConfig( ) } override def assignBreakpoint( - topology: Array[WorkerLayer], - states: mutable.AnyRefMap[ActorRef, WorkerState.Value], breakpoint: GlobalBreakpoint - )(implicit timeout: Timeout, ec: ExecutionContext): Unit = { - breakpoint.partition(topology(0).layer.filter(states(_) != WorkerState.Completed)) + ): Array[ActorVirtualIdentity] = { + topology.layers(0).identifiers } } diff --git a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/sink/SimpleSinkOpExecConfig.scala b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/sink/SimpleSinkOpExecConfig.scala index 183c65e76df..2bbaa4de432 100644 --- a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/sink/SimpleSinkOpExecConfig.scala +++ b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/sink/SimpleSinkOpExecConfig.scala @@ -8,6 +8,7 @@ import edu.uci.ics.amber.engine.architecture.deploysemantics.deploymentfilter.Fo import edu.uci.ics.amber.engine.architecture.deploysemantics.deploystrategy.RandomDeployment import edu.uci.ics.amber.engine.architecture.deploysemantics.layer.WorkerLayer import edu.uci.ics.amber.engine.architecture.worker.WorkerState +import edu.uci.ics.amber.engine.common.ambertag.neo.VirtualIdentity.ActorVirtualIdentity import edu.uci.ics.amber.engine.common.ambertag.{LayerTag, OperatorIdentifier} import edu.uci.ics.amber.engine.operators.SinkOpExecConfig @@ -30,10 +31,8 @@ class SimpleSinkOpExecConfig(tag: OperatorIdentifier) extends SinkOpExecConfig(t ) override def assignBreakpoint( - topology: Array[WorkerLayer], - states: mutable.AnyRefMap[ActorRef, WorkerState.Value], breakpoint: GlobalBreakpoint - )(implicit timeout: Timeout, ec: ExecutionContext): Unit = { - breakpoint.partition(topology(0).layer.filter(states(_) != WorkerState.Completed)) + ): Array[ActorVirtualIdentity] = { + topology.layers(0).identifiers } } diff --git a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/source/mysql/MysqlSourceOpExecConfig.scala b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/source/mysql/MysqlSourceOpExecConfig.scala index 9e9a3c927ca..c72c2328af1 100644 --- a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/source/mysql/MysqlSourceOpExecConfig.scala +++ b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/source/mysql/MysqlSourceOpExecConfig.scala @@ -7,6 +7,7 @@ import edu.uci.ics.amber.engine.architecture.deploysemantics.deploymentfilter.Us import edu.uci.ics.amber.engine.architecture.deploysemantics.deploystrategy.OneOnEach import edu.uci.ics.amber.engine.architecture.deploysemantics.layer.WorkerLayer import edu.uci.ics.amber.engine.architecture.worker.WorkerState +import edu.uci.ics.amber.engine.common.ambertag.neo.VirtualIdentity.ActorVirtualIdentity import edu.uci.ics.amber.engine.common.ambertag.{LayerTag, OperatorIdentifier} import edu.uci.ics.amber.engine.operators.OpExecConfig import edu.uci.ics.texera.workflow.common.operators.source.SourceOperatorExecutor @@ -34,13 +35,9 @@ class MysqlSourceOpExecConfig( Map() ) } - override def assignBreakpoint( - topology: Array[WorkerLayer], - states: mutable.AnyRefMap[ActorRef, WorkerState.Value], breakpoint: GlobalBreakpoint - )(implicit timeout: Timeout, ec: ExecutionContext): Unit = { - breakpoint.partition(topology(0).layer.filter(states(_) != WorkerState.Completed)) + ): Array[ActorVirtualIdentity] = { + topology.layers(0).identifiers } - } diff --git a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/visualization/pieChart/PieChartOpExecConfig.scala b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/visualization/pieChart/PieChartOpExecConfig.scala index 8c85d7fab97..f1638851092 100644 --- a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/visualization/pieChart/PieChartOpExecConfig.scala +++ b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/visualization/pieChart/PieChartOpExecConfig.scala @@ -13,6 +13,7 @@ import edu.uci.ics.amber.engine.architecture.deploysemantics.layer.WorkerLayer import edu.uci.ics.amber.engine.architecture.linksemantics.HashBasedShuffle import edu.uci.ics.amber.engine.architecture.worker.WorkerState import edu.uci.ics.amber.engine.common.Constants +import edu.uci.ics.amber.engine.common.ambertag.neo.VirtualIdentity.ActorVirtualIdentity import edu.uci.ics.amber.engine.common.ambertag.{LayerTag, OperatorIdentifier} import edu.uci.ics.amber.engine.operators.OpExecConfig import edu.uci.ics.texera.workflow.common.tuple.Tuple @@ -62,11 +63,9 @@ class PieChartOpExecConfig( } override def assignBreakpoint( - topology: Array[WorkerLayer], - states: mutable.AnyRefMap[ActorRef, WorkerState.Value], breakpoint: GlobalBreakpoint - )(implicit timeout: Timeout, ec: ExecutionContext): Unit = { - breakpoint.partition(topology(0).layer.filter(states(_) != WorkerState.Completed)) + ): Array[ActorVirtualIdentity] = { + topology.layers(0).identifiers } } diff --git a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/visualization/wordCloud/WordCloudOpExecConfig.scala b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/visualization/wordCloud/WordCloudOpExecConfig.scala index 4240c59e842..ae81003753d 100644 --- a/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/visualization/wordCloud/WordCloudOpExecConfig.scala +++ b/core/amber/src/main/scala/edu/uci/ics/texera/workflow/operators/visualization/wordCloud/WordCloudOpExecConfig.scala @@ -13,6 +13,7 @@ import edu.uci.ics.amber.engine.architecture.deploysemantics.layer.WorkerLayer import edu.uci.ics.amber.engine.architecture.linksemantics.HashBasedShuffle import edu.uci.ics.amber.engine.architecture.worker.WorkerState import edu.uci.ics.amber.engine.common.Constants +import edu.uci.ics.amber.engine.common.ambertag.neo.VirtualIdentity.ActorVirtualIdentity import edu.uci.ics.amber.engine.common.ambertag.{LayerTag, OperatorIdentifier} import edu.uci.ics.amber.engine.operators.OpExecConfig import edu.uci.ics.texera.workflow.common.tuple.Tuple @@ -60,11 +61,9 @@ class WordCloudOpExecConfig( } override def assignBreakpoint( - topology: Array[WorkerLayer], - states: mutable.AnyRefMap[ActorRef, WorkerState.Value], breakpoint: GlobalBreakpoint - )(implicit timeout: Timeout, ec: ExecutionContext): Unit = { - breakpoint.partition(topology(0).layer.filter(states(_) != WorkerState.Completed)) + ): Array[ActorVirtualIdentity] = { + topology.layers(0).identifiers } } diff --git a/core/amber/src/test/scala/edu/uci/ics/amber/engine/architecture/breakpoint/ExceptionBreakpointSpec.scala b/core/amber/src/test/scala/edu/uci/ics/amber/engine/architecture/breakpoint/ExceptionBreakpointSpec.scala index 190d949e721..aa22dd6bb0f 100644 --- a/core/amber/src/test/scala/edu/uci/ics/amber/engine/architecture/breakpoint/ExceptionBreakpointSpec.scala +++ b/core/amber/src/test/scala/edu/uci/ics/amber/engine/architecture/breakpoint/ExceptionBreakpointSpec.scala @@ -1,10 +1,6 @@ package edu.uci.ics.amber.engine.architecture.breakpoint import edu.uci.ics.amber.clustering.SingleNodeListener -import edu.uci.ics.amber.engine.architecture.breakpoint.globalbreakpoint.{ - ConditionalGlobalBreakpoint, - CountGlobalBreakpoint -} import edu.uci.ics.amber.engine.architecture.controller.{Controller, ControllerState} import edu.uci.ics.amber.engine.common.AdvancedMessageSending import edu.uci.ics.amber.engine.common.ambermessage.ControlMessage.{ diff --git a/core/amber/src/test/scala/edu/uci/ics/amber/engine/architecture/controller/ControllerSpec.scala b/core/amber/src/test/scala/edu/uci/ics/amber/engine/architecture/controller/ControllerSpec.scala index 595be27009b..c8a7a0388b1 100644 --- a/core/amber/src/test/scala/edu/uci/ics/amber/engine/architecture/controller/ControllerSpec.scala +++ b/core/amber/src/test/scala/edu/uci/ics/amber/engine/architecture/controller/ControllerSpec.scala @@ -1,10 +1,6 @@ package edu.uci.ics.amber.engine.architecture.controller import edu.uci.ics.amber.clustering.SingleNodeListener -import edu.uci.ics.amber.engine.architecture.breakpoint.globalbreakpoint.{ - ConditionalGlobalBreakpoint, - CountGlobalBreakpoint -} import edu.uci.ics.amber.engine.common.ambermessage.ControlMessage.{ Ack, ModifyLogic, diff --git a/core/amber/src/test/scala/edu/uci/ics/amber/engine/e2e/DataProcessingSpec.scala b/core/amber/src/test/scala/edu/uci/ics/amber/engine/e2e/DataProcessingSpec.scala index 99e8f46a86e..603f4264727 100644 --- a/core/amber/src/test/scala/edu/uci/ics/amber/engine/e2e/DataProcessingSpec.scala +++ b/core/amber/src/test/scala/edu/uci/ics/amber/engine/e2e/DataProcessingSpec.scala @@ -1,10 +1,6 @@ package edu.uci.ics.amber.engine.e2e import edu.uci.ics.amber.clustering.SingleNodeListener -import edu.uci.ics.amber.engine.architecture.breakpoint.globalbreakpoint.{ - ConditionalGlobalBreakpoint, - CountGlobalBreakpoint -} import edu.uci.ics.amber.engine.common.ambermessage.ControlMessage.{ Ack, ModifyLogic,