【发布时间】:2014-06-25 08:52:41
【问题描述】:
我有一个 Python 脚本,我在其中绑定多个(例如 5 个)ZMQ 接收器套接字,如下所示:
receiver_1 = context.socket(zmq.PULL)
receiver_1.bind("tcp://*:5555")
...
receiver_5 = context.socket(zmq.PULL)
receiver_5.bind("tcp://*:5559")
receivers = [receiver_1, ..., receiver_5]
然后我启动一些 Google Compute Engine 实例,它们连接相应的发送方套接字。
我想从这些套接字中并行提取,因此我尝试使用多处理池来实现。代码如下所示:
def recv_result(i):
result_str = receivers[i].recv()
return cPickle.loads(result_str)
pool = multiprocessing.Pool()
while True:
results = pool.map(recv_result, [i for i in range(len(receivers))])
# break when all results have been received
...
我在运行脚本时遇到的错误如下所示:
Traceback (most recent call last):
...
File ...
results = pool.map(recv_result, [i for i in range(len(receivers))])
File "/usr/lib/python2.7/multiprocessing/pool.py", line 227, in map
return self.map_async(func, iterable, chunksize).get()
File "/usr/lib/python2.7/multiprocessing/pool.py", line 528, in get
raise self._value
zmq.error.ZMQError: Interrupted system call
我也尝试过使用 multiprocessing.Process 实现相同的功能,但我得到了基本相同的错误,尽管方式更加混乱。
我想要做的是更有效地接收来自我的 GCE 实例的所有结果,因为我发现这是我脚本中的瓶颈(在我当前的工作实现中,我只有一个接收器套接字,并且它串行接收来自所有 GCE 实例的结果)。如果有人能指出我当前代码中的错误,或者任何关于实现目标的更好方法的建议,将不胜感激!
【问题讨论】:
标签: python multiprocessing zeromq pyzmq