【更好的理解ThingsBoard规则引擎的】深入理解Akka Actor模型
Carl Hewitt 在1973年对Actor模型进行了如下定义:"Actor模型是一个把'Actor'作为并发计算的通用原语". Actor是异步驱动,可以并行和分布式部署及运行的最小颗粒。也就是说,它可以被分配,分布,调度到不同的CPU,不同的节点,乃至不同的时间片上运行,而不影响最终的结果。因此Actor在空间(分布式)和时间(异步驱动)上解耦的。而Akka是Lightbend(前身是Typesafe)公司在JVM上的Actor模型的实现。我们在了解actor模型之前,首先来了解actor模型主要是为了解决什么样的问题。
Why modern systems need a new programming model
在akka系统的官网上主要介绍了现代并发编程模型所遇到的问题,里面主要提到了三个点
1) 在面向对象的语言中一个显著的特点是封装,然后通过对象提供的一些方法来操作其状态,但是共享内存的模型下,多线程对共享对象的并发访问会造成并发安全问题。一般会采用加锁的方式去解决

加锁会带来一些问题
加锁的开销很大,线程上下文切换的开销大
加锁导致线程block,无法去执行其他的工作,被block无法执行的线程,其实也是占据了一种系统资源
加锁在编程语言层面无法防止隐藏的死锁问题
2)我们知道Java中并发模型是通过共享内存来实现。而cpu中会利用局部cache来加速主存的访问,为了解决多线程间缓存不一致的问题,在java中一般会通过使用volatile或者Atmoic来标记变量,通过Jmm的happens before机制来保障多线程间共享变量的可见性。因此从某种意义上来说是没有共享内存的,而是通过cpu将cache line的数据刷新到主存的方式来实现可见。
因此与其去通过标记共享变量或者加锁的方式,依赖cpu缓存更新,倒不如每个并发实例之间只保存local的变量,而在不同的实例之间通过message来传递。
3)call stack的问题
当我们编程模型异步化之后,还有一个比较大的问题是调用栈转移的问题,如下图中主线程提交了一个异步任务到队列中,Worker thread 从队列提取任务执行,调用栈就变成了workthread发起的,当任务出现异常时,处理和排查就变得困难。

How the Actor Model Meets the Needs of Modern Distributed Systems
那么akka 的actor的模型是怎样处理这些问题的,actor模型中的抽象主体变为了actor,
actor之间可以互相发送message。
actor在收到message之后会将其存入其绑定的Mailbox中。
Actor从Mailbox中提取消息,执行内部方法,修改内部状态。
继续给其他actor发送message。
可以看到下图,actor内部的执行流程是顺序的,同一时刻只有一个message在进行处理,也就是actor的内部逻辑可以实现无锁化的编程。actor和线程数解耦,可以创建很多actor绑定一个线程池来进行处理,no lock,no block的方式能减少资源开销,并提升并发的性能

actor编程样例
下面简单来看一个actor的样例
依赖
<dependency> <groupId>com.typesafe.akka</groupId> <artifactId>akka-actor_2.11</artifactId> <version>2.4.20</version> </dependency>
Main
public static void main(String[] args) throws InterruptedException {
final ActorSystem actorSystem = ActorSystem.create("actor-system");
final ActorRef actorRef = actorSystem.actorOf(Props.create(BankActor.class), "bank-actor");
CountDownLatch addCount = new CountDownLatch(20);
CountDownLatch minusCount = new CountDownLatch(10);
Thread addCountT = new Thread(new Runnable() {
@Override
public void run() {
while (addCount.getCount() > 0) {
actorRef.tell(Command.ADD, null);
addCount.countDown();
}
}
});
Thread minusCountT = new Thread(new Runnable() {
@Override
public void run() {
while (minusCount.getCount() > 0) {
actorRef.tell(Command.MINUS, null);
minusCount.countDown();
}
}
});
minusCountT.start();
addCountT.start();
minusCount.await();
addCount.await();
Future<Object> count = Patterns.ask(actorRef, Command.GET, 1000);
count.onComplete(
new OnComplete<Object>() {
@Override
public void onComplete(Throwable failure, Object success) throws Throwable {
if (failure != null) {
failure.printStackTrace();
} else {
log.info("Get result from " + success);
}
}
},
Executors.directExecutionContext());
actorSystem.shutdown();
}创建actor
通过actorRef和actor并发交互
获取actor最后的状态
actor
public class BankActor extends UntypedActor {
private static final Logger log = LoggerFactory.getLogger(BankActor.class);
private int count;
@Override
public void preStart() throws Exception, Exception {
super.preStart();
count = 0;
}
@Override
public void onReceive(Object message) throws Throwable {
// 可以使用枚举或者动态代理类来实现方法调用
if (message instanceof Command) {
Command cmd = (Command) message;
switch (cmd) {
case ADD:
log.info("Add 1 from {} to {}", count, ++count);
break;
case MINUS:
log.info("Minus 1 from {} to {}", count, --count);
break;
case GET:
log.info("Return current count " + getSender());
getSender().tell(count, this.getSelf());
break;
default:
log.warn("U原创不易,完成人机校验,阅读全文