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ç‰ï¼‰æ�¥å�‘é€�。