12.10 定义一个Actor任务¶

问题¶

你想定义跟actor模�中类似“actors�角色的任务

解决方案¶

actor模�是一�最��的也是最简�的并行和分布�计算解决方案。 事实上,它天生的简�性是它如此�欢迎的��原因之一。 简��讲,一个actor就是一个并�执行的任务,�是简�的执行��给它的消�任务。 �应这些消�时,它�能还会给其他actor��更进一步的消�。 actor之间的通信是��和异步的。因此,消���者�知�消�是什么时候被��, 也�会接收到一个消�已被处�的回应或通知。

结�使用一个线程和一个队列�以很容易的定义actor,例如:

from queue import Queue
from threading import Thread, Event

# Sentinel used for shutdown
class ActorExit(Exception):
    pass

class Actor:
    def __init__(self):
        self._mailbox = Queue()

    def send(self, msg):
        '''
        Send a message to the actor
        '''
        self._mailbox.put(msg)

    def recv(self):
        '''
        Receive an incoming message
        '''
        msg = self._mailbox.get()
        if msg is ActorExit:
            raise ActorExit()
        return msg

    def close(self):
        '''
        Close the actor, thus shutting it down
        '''
        self.send(ActorExit)

    def start(self):
        '''
        Start concurrent execution
        '''
        self._terminated = Event()
        t = Thread(target=self._bootstrap)

        t.daemon = True
        t.start()

    def _bootstrap(self):
        try:
            self.run()
        except ActorExit:
            pass
        finally:
            self._terminated.set()

    def join(self):
        self._terminated.wait()

    def run(self):
        '''
        Run method to be implemented by the user
        '''
        while True:
            msg = self.recv()

# Sample ActorTask
class PrintActor(Actor):
    def run(self):
        while True:
            msg = self.recv()
            print('Got:', msg)

# Sample use
p = PrintActor()
p.start()
p.send('Hello')
p.send('World')
p.close()
p.join()

这个例�中,你使用actor实例的 send() 方法��消�给它们。 其机制是,这个方法会将消�放入一个队里中, 然�将其转交给处�被接�消�的一个内部线程。 close() 方法通过在队列中放入一个特殊的哨兵值(ActorExit)�关闭这个actor。 用户�以通过继承Actor并定义实现自己处�逻辑run()方法�定义新的actor。 ActorExit 异常的使用就是用户自定义代��以在需�的时候��获终止请求 (异常被get()方法抛出并传播出去)。

如果你放宽对于�步和异步消���的�求, 类actor对象还�以通过生�器�简化定义。例如:

def print_actor():
    while True:

        try:
            msg = yield      # Get a message
            print('Got:', msg)
        except GeneratorExit:
            print('Actor terminating')

# Sample use
p = print_actor()
next(p)     # Advance to the yield (ready to receive)
p.send('Hello')
p.send('World')
p.close()

讨论¶

actor模�的魅力就在于它的简�性。 实际上,这里仅仅�有一个核心�作 send() . 甚至,对于在基于actor系统中的“消��的泛化概念�以已多�方�被扩展。 例如,你�以以元组形�传递标签消�,让actor执行��的�作,如下:

class TaggedActor(Actor):
    def run(self):
        while True:
             tag, *payload = self.recv()
             getattr(self,'do_'+tag)(*payload)

    # Methods correponding to different message tags
    def do_A(self, x):
        print('Running A', x)

    def do_B(self, x, y):
        print('Running B', x, y)

# Example
a = TaggedActor()
a.start()
a.send(('A', 1))      # Invokes do_A(1)
a.send(('B', 2, 3))   # Invokes do_B(2,3)
a.close()
a.join()

作为�外一个例�,下�的actor�许在一个工作者中�行任�的函数, 并且通过一个特殊的Result对象返回结果:

from threading import Event
class Result:
    def __init__(self):
        self._evt = Event()
        self._result = None

    def set_result(self, value):
        self._result = value

        self._evt.set()

    def result(self):
        self._evt.wait()
        return self._result

class Worker(Actor):
    def submit(self, func, *args, **kwargs):
        r = Result()
        self.send((func, args, kwargs, r))
        return r

    def run(self):
        while True:
            func, args, kwargs, r = self.recv()
            r.set_result(func(*args, **kwargs))

# Example use
worker = Worker()
worker.start()
r = worker.submit(pow, 2, 3)
worker.close()
worker.join()
print(r.result())

最�,“���一个任务消�的概念�以被扩展到多进程甚至是大型分布�系统中去。 例如,一个类actor对象的 send() 方法�以被编程让它能在一个套接字连接上传输数� 或通过�些消�中间件(比如AMQP�ZMQ等)���。