JAVA游戏

【更好的理解ThingsBoard规则引擎的】深入理解Akka Actor模型

CarlHewitt在1973年对Actor模型进行了如下定义:"Actor模型是一个把'Actor'作为并发计算的通用原语".Actor是异步驱动,可以并行和分布式部署及运行的最小颗粒。也就是说,它可以被分配,分布,调度到不同的CPU,不同的节点,乃至不同的时间片上运行,而不影响最终的结果。因此Actor在空间(分布式)和时间(异步驱动)上解耦的。而Akka是

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) 在面向对象的语言中一个显著的特点是封装,然后通过对象提供的一些方法来操作其状态,但是共享内存的模型下,多线程对共享对象的并发访问会造成并发安全问题。一般会采用加锁的方式去解决

473fa06b23b1a5c4ab25d440acba9e81.png

加锁会带来一些问题

  • 加锁的开销很大,线程上下文切换的开销大

  • 加锁导致线程block,无法去执行其他的工作,被block无法执行的线程,其实也是占据了一种系统资源

  • 加锁在编程语言层面无法防止隐藏的死锁问题


2)我们知道Java中并发模型是通过共享内存来实现。而cpu中会利用局部cache来加速主存的访问,为了解决多线程间缓存不一致的问题,在java中一般会通过使用volatile或者Atmoic来标记变量,通过Jmm的happens before机制来保障多线程间共享变量的可见性。因此从某种意义上来说是没有共享内存的,而是通过cpu将cache line的数据刷新到主存的方式来实现可见。
因此与其去通过标记共享变量或者加锁的方式,依赖cpu缓存更新,倒不如每个并发实例之间只保存local的变量,而在不同的实例之间通过message来传递。

3)call stack的问题
当我们编程模型异步化之后,还有一个比较大的问题是调用栈转移的问题,如下图中主线程提交了一个异步任务到队列中,Worker thread 从队列提取任务执行,调用栈就变成了workthread发起的,当任务出现异常时,处理和排查就变得困难。

73d3450994de001f62528d0553c8695a.png

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的方式能减少资源开销,并提升并发的性能

65c076f57af0592d69f0bc680d1b4976.png

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();
    }


  1. 创建actor

  2. 通过actorRef和actor并发交互

  3. 获取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

原创不易,完成人机校验,阅读全文

相关推荐