【问题标题】:Akka persistent actor - functional approachAkka 持久化参与者 - 函数式方法
【发布时间】:2020-03-20 04:32:40
【问题描述】:

我一直在尝试使用函数式方法来实现持久性参与者——我的意思是根本没有变量。但是我遇到了麻烦 :) 下面的代码示例不能正常工作,因为处理程序参数在处理程序之间不共享(receiveCommand/receiveRecover)。两者都从零开始,然后相互覆盖 - 重播某些事件后,命令处理程序仍将位于起点。

这个实现的另一个问题是在同一个地方处理命令和事件

以实用的方式完成它是一种好习惯吗?

class Item(sku: String) extends PersistentActor with ActorLogging {
    import Item._
    override def persistenceId: String = sku
    override def receiveCommand: Receive = handler(0, 0)
    override def receiveRecover: Receive = handler(0, 0)

    def handler(quantity: Int, booked: Int): Receive = {
      case Increase(q) =>
        val event = StockChanged(sku, q)
        persist(event)(e => context.become(handler(quantity + e.quantity, booked)))
      case Decrease(q) =>
        val event = StockChanged(sku, -q)
        persist(event)(e => context.become(handler(quantity + e.quantity, booked)))
      case StockChanged(_, q) => {
        context.become(handler(quantity + q, booked))
      }
    }
}

【问题讨论】:

  • 不太清楚你的问题是什么。请提出一个特定的问题,指出一个特定的问题,该问题可以被没有你的整个代码库的其他人复制。
  • 我的问题很简单,如何在不使用变量(保持其状态)的情况下实现持久化actor
  • 所以actor是无国籍的?让你的演员无国籍有什么问题?
  • 不,不是。我想将其状态实现为函数
  • 演员的状态是否由您在其中的handle 函数表示?

标签: scala functional-programming akka actor


【解决方案1】:

我找到了答案。新的 akka 持久性 2.6 只是以功能方式工作:)

https://doc.akka.io/docs/akka/current/typed/persistence.html#event-sourcing

【讨论】:

    【解决方案2】:

    对于那些需要使用Classic Persistence 的人来说,我能想出的最佳方法不需要演员范围的状态,但你仍然需要一个本地的var。然后,您可以将本地 var 封装在恢复处理程序的闭包中。

    扩展和清理问题中不完整的例子:

    package item_example
    
    import akka.actor.ActorLogging
    import akka.persistence.PersistentActor
    
    object ItemExample {
    
      case class Increase(quantity: Int)
    
      case class Decrease(quantity: Int)
    
      case class StockChanged(id: String, quantity: Int)
    
      class Item(sku: String) extends PersistentActor with ActorLogging {
        override def persistenceId: String = sku
    
        override def receiveCommand: Receive = listening()
    
        private def listening(quantity: Int = 0): Receive = {
          case Increase(q) =>
            persist(StockChanged(sku, q)) { e: StockChanged =>
              context.become(listening(quantity + e.quantity))
            }
          case Decrease(q) =>
            persist(StockChanged(sku, -q)) { e: StockChanged =>
              context.become(listening(quantity + e.quantity))
            }
        }
    
        override def receiveRecover: Receive = recovering()
    
        private def recovering(): Receive = {
          /*
          This is the local var that will be part of the recovery handler's closure.
          The semicolon is needed only so that the below code within parentheses is interpreted as
          a partial function and not a code block parameter for the integer 0.
          You can also just return a new PartialFunction[Any, Unit] but I find this to be nicer.
           */
          var totalQuantity: Int = 0;
          {
            case StockChanged(_, quantity) =>
              totalQuantity += quantity
              context.become(listening(totalQuantity))
          }
        }
      }
    }
    

    恢复处理程序将覆盖每个恢复事件后的状态。这没关系,因为recovery is guaranteed to finish before the actor starts consuming events from the mailbox

    上述方法有效,因为在恢复处理程序中调用context::become 不会修改恢复处理程序!它只修改邮箱处理程序。当我们处理重播的事件时,我们只需在每个事件之后替换邮箱处理程序。 对于所有重放的事件,恢复处理程序将被原样调用。

    旁注:区分发送给参与者的命令和从持久存储重放的事件很重要。该问题定义了两个命令:IncreaseDecrease。它还定义了一个事件:StockChanged。但是,receiveCommand 方法处理所有这三个。这是一种代码气味。 receiveCommand 应该处理命令,而receiveRecover 应该处理事件。我已经在上面的解决方案中清理了这个问题。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2014-10-19
      • 2017-11-30
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多