Skip to main content

用于长时间运行的多处理事件循环的库。

项目描述

# mp_event_loop

库,用于长时间运行的多处理事件循环。该库提供了一个 EventLoop,它将在
单独的进程中运行事件。该库的目的是在 GUI 在主
线程中运行时管理长时间运行的进程。当 GUI 运行时,任务可以连续卸载到另一个进程。

**现在可以与 async/await 一起使用!** 请参见下面的示例。

EventLoop 带有几个用于管理单独进程的实用程序。

* is_running() - 如果单独的进程正在运行则返回
* start() - 启动事件循环
* run(events=None, output_handlers=None) - 添加输出处理程序并运行事件
* run_until_complete(events=None, output_handlers=None) - 运行给定事件并等待它们完成
* wait() - 等待当前事件完成
* stop() - 停止事件循环
* close() - 关闭事件循环
* _\_enter_\_ 和 _\_exit_\_ - 用作上下文管理器并允许使用 `with` 语句

EventLoop 还附带一些实用程序来添加要处理的命令和处理结果的方法。

* add_event(target, args, kwargs, ...) - 添加要在单独进程中执行的事件。
* add_output_handler(function) - 在事件执行后接收事件的函数。

这些功能将在下面详细解释


##快速入门

pip install mp_event_loop

```python
import mp_event_loop


def add_vals(value, value2=1):
return value + value2


results = []

def save_results(event):
results.append(event.results)


with mp_event_loop.get_event_loop(save_results):
mp_event_loop.add_event( add_vals, 2)
mp_event_loop.add_event(add_vals, 3, 4)
mp_event_loop.add_event(add_vals, 5, value2=6)
mp_event_loop.add_event(add_vals, args=(7,), kwargs={'value2': 8})

# with context manager 等待事件和事件结果完成,直到它退出
assert results == [3, 7, 11, 15]
summed = sum(results)
assert summed == 36
print("Results summed:", summed)
# 结果总和:
36```

替代方法
```python
import mp_event_loop

def add_one(value):
return value + 1


results = []

def save_results(event):
results.append(event.results)

mp_event_loop.run([{'target': add_one, 'args ': (1,)},
{'target': add_one, 'args': (2,)},
{'target': add_one, 'args': (3,)},
], save_results)

mp_event_loop.add_event( add_one, 4)

mp_event_loop.wait()

print("Results summed:", sum(results))
# Results summed: 14
```
## Async /

Await 您必须在 AsyncManager 中注册您的协程才能在单独的过程。协程不是
picklable,所以它们必须在模块级别注册。在 unpickling (_\_setstate_\_) 期间,将
使用 AsyncManager 通过注册名称检索协程。


```python
from mp_event_loop import AsyncEventLoop, AsyncManager


async def print_test(value, name):
print("Print", name)
return value


async def yield_range(value, name):
print("Yield", name)
for i in range (value):
yield name + " " + str(i)


# 不再需要注册
# AsyncManager.register('print_test', print_test)
# AsyncManager.register('yield_range', yield_range)


if __name__ == '__main__' :



results.append(event.results)

与 AsyncEventLoop(output_handlers=save_results) 作为循环:
loop.async_event(print_test, 1, "hello")
loop.async_event(print_test, 2, "hi")
loop.async_event('print_test', 3, "oi") # 也可以使用注册名
loop.async_event(yield_range, 5, 'first')
loop.async_event(yield_range, 5, 'second')

print(results)
# [1, 2, 3, '第一个0','第一个1','第二个0','第一个2','第二个1','第一个3','第二个2','第一个4','第二个3','第二个4']

` ``

## 工作原理
EventLoop 通过创建进程和线程来工作。该过程采取
来自队列的事件并运行函数。一旦事件完成,事件就会被放入结果队列/消费者队列中。
线程从结果队列中获取事件并将其传递给事件循环中的所有输出处理程序。
如果其中一个 output_handlers 返回 True,则事件将停止传播到其他 output_handlers。

由于队列中的锁定机制和进程之间的消息传递,这将很慢。您可能
只会将其用于并发性。这对于线程可能影响性能的非 IO 并发很有用。

我创建了这个库作为一个测试来了解多处理是如何工作的。我正在尝试使用多处理
tcp 通信和解析数据。我希望解析发生在一个单独的进程中,但我想
在允许 GUI 运行的线程中访问解析的数据。


## Object Persistence
这个库的总体目标是一个长时间运行的进程,任务可以轻松运行。随之而来的是对象
持久性。您可以轻松地将对象发送到单独的进程并在其上运行任务。困难的部分是
从该对象中获取值。多处理库已经作为通过 Manager
类执行此操作的好方法。但是,我发现这种方法很难与 PySide 一起使用。

```python
# 示例问题
import mp_event_loop


class ABC(object):
def __init__(self, a, b, c):
self.a = a
self.b = b
self.c = c

def calc_c(self):
self.c = self.a + self.b

def __repr__(self):
return "ABC(a=%d, b=% d, c=%d)" % (self.a, self.b, self.c)


a = ABC(1, 2, 0)

with mp_event_loop.get_event_loop(has_results=False) 作为循环:
loop.add_event(print, "=====", a, "=====") # 正确打印 a "ABC(a=1, b=2, c=0)"
loop.add_event(a.calc_c) # 计算,但没有将结果发回给我们
# 另一个进程具有正确的 c 值,但从未将其传递回正确的对象

# 我们将这个 a 的值向下传递,即 1、2、0。我们从未得到结果“a.calc_c()”
loop.add_event(print, a) # 打印主线程传递给进程的 'a' "ABC(a=1, b=2, c=0)"
```

这个库试图分两步解决这个问题方法。

我标记为缓存的第一种方式。CacheEvent 通过在字典中保存对对象的引用来解决这个问题。
对象被一次腌制到一个单独的进程并保存。每个后续的 CacheEvent 都使用一个键 (id) 来
检索缓存对象并使用该缓存对象在另一个进程中运行

。第二种方式是通过代理,其中一个对象仅存在于另一个进程中。我的代理方法的
主要过程是保持对简单对象的引用。当该对象被腌制时,它会在另一个对象中创建一个对象
过程。它创建该对象一次,然后使用缓存来保存对该对象的引用。该对象仅存在于
其他进程中。您可以通过调用主进程中的函数来控制对象,该函数实际上只是发送一个
事件以在另一个进程上运行。如果你想在主进程上暴露一个属性值或 getter 方法值,
那么你需要定义 `PROPERTIES` 和 `GETTERS`。注意:这只是我发现使用 PySide 并
在另一个进程中动态创建小部件的解决方案。我计划拥有另一个名为 qt_concurrency 的库来处理
长时间运行的多处理和 Qt 之间的细节。


### 缓存事件
上面的例子表明主进程中的 `a` 永远不会改变。`a.calc_c()` 正在运行其他进程,但
其他进程不保存对象 `a`。在另一个进程的范围内,带有 `a.calc_c()` 的 `a` 会死掉并消失。
主进程“a”从未改变,仍然具有 a=1、b=2、c=0 的值。

为了解决这个问题,我创建了一些缓存事件。这些事件使用 id 保存对象。另一个进程被赋予
object_id 并使用它来获取正确的对象。

```python
a = ABC(1, 2, 0)

with mp_event_loop.get_event_loop(has_results=False) 作为循环:
loop.add_event(print, "=====", a, "=====") # 打印正确的 "ABC(a=1, b=2, c=0)"

# 你可以手动缓存一个对象
# loop.cache_object(a)
# 只有对象参数需要'loop.add_event(my_func, a, cache=True)' 只有my_func被注册和缓存。
# 如果 'a' 被缓存,它将通过使用 a 的 object_id 来使用它,但它不会注册和缓存 'a' 对象。

# 或者传入 cache=True 或 loop.add_cache_event
loop.add_event(a.calc_c, cache=True) # 跟踪 'a' 的 id 并将 'a' 保存在单独的进程中

# 不要将 a down 传递给其他过程。而是为“a”传递一个 object_id。另一个进程将使用
#object_id 来获取存储的“a”对象并使用该对象。
loop.add_event(print, a, cache=True) # 正确打印缓存的 'a' 对象 "ABC(a=1, b=2,


虽然这现在像我们希望的那样工作,但仍然存在问题。主进程中的 `a` 没有正确的
值。您可以通过如下所示的繁琐事件处理来解决此问题。

```python
import mp_event_loop

a = ABC(1, 2, 0)


def save_object(event):
if event.event_key == 'a':
aa = event.results.a
ab = event.results.b
ac = event. results.c
return True


with mp_event_loop.get_event_loop(output_handlers=save_object) as loop:
loop.add_event(print, "=====", a, "=====") # 正确打印一个 "ABC(a= 1, b=2, c=0)"
loop.add_event(a.calc_c) # 计算并通过output_handler发回事件
# 但是我们无法确定返回给我们的事件/对象(在 save_object 中)

# 尝试使用 save_result_object 取回对象
loop.add_event(a.calc_c, event_key='a') # 我们可以使用 event_key识别返回的事件 (in save_object)
print("Main process", a) # 多进程还没有发生。没有立竿见影的效果。

print("事件循环后的主进程", a) # 循环现在已经等待所有任务完成并且结果正确。
```

### 代理对象
这是迄今为止这个库中最具挑战性的部分。对于我的应用程序,我不关心立竿见影的结果。我
关心在允许 GUI 存在的单独进程中完成的并发和长时间运行的任务。

下面是代理的示例代码。代理对象只能访问

```python
import mp_event_loop

class Point(object):
def __init__(self, x=0, y=0, z=0):
self.x = x
self.y = y
self._z = z

def get_z(self):
return self._z

def set_z(self, z):
self._z = z

def move(self, x, y, z):
self.x = x
self.y = y
self._z = z

class MpPoint(mp_event_loop.Proxy):
PROXY_CLASS = Point
PROPERTIES = ['x', 'y']
GETTERS = ['get_z']

def __init__(self, x=0, y=0, z=0, loop=无):
super().__init__(loop=loop)
self.x = x
self.y = y
self._z = z



with mp_event_loop.EventLoop() as loop:
p = MpPoint(loop=loop)
p.set_z(3) # 在其他进程中运行
assert p._z != 3
assert p.get_z() != 3

# 将对象发送到单独的进程并运行 move。
# 单独的进程具有属性 x 和 y 并使用 get_z 和 set_z 来获得正确的 z 值。
p.move(1, 2, 7)
断言 px == 0
断言 py == 0
断言 p._z == 3
断言 p.get_z() == 3

断言 px == 1
断言 py == 2
断言 p.get_z () == 7
```

## 事件
事件只是接受一个函数和一些参数,并在一个单独的进程中执行它们。从理论上讲,您可以
为特定的事情制作自己的事件。

```python
import mp_event_loop

class MyEvent(mp_event_loop.Event):
def __init__(self, data, **kwargs):
super().__init__(target=None, args=(data,), **kwargs)

def run( self):
data = self.args[0]

# 运行一些计算
value = list(range(data))

return value

# def exec_(self):
# """获取命令并运行它"""
# # 获取命令运行
# self.results = None
# self.error = None
# if callable(self.target):
# # 运行命令
# try:
# self.results = self.run()
# except Exception as err:
# self.error = err
# elif self.target is not None:
# self .error = ValueError("Invalid target (%s) given! Type %s" % (repr(self.target), str(type(self.target))))


def print_results(event):
print(event.results)


loop = mp_event_loop.EventLoop(output_handlers=print_results)

loop.start()

loop.add_event(MyEvent(10))

loop.stop()

# 在某个点 [0, 1, 2, 3, 4, 5, 6, 7, 8, 9] 应该打印
```

将函数传递给事件要容易得多,但您可能会发现这很有用。


## 输出处理程序

EventLoop 包含一个输出处理程序列表。输出处理程序只是一个接收事件的简单函数。
Event 对象将具有 results 属性,其中包含事件执行的结果(
从目标函数返回的结果)。

```python
import mp_event_loop

class MyEvent(mp_event_loop.Event):
def __init__(self, data, **kwargs):
super().__init__(target=None, args=(data,), **kwargs)

def run( self):
data = self.args[0]

# 运行一些计算
value = list(range(data))

返回值


def print_my_event(event):
if isinstance(event, MyEvent):
print('My Event', event.results)
return True # 停止运行其他 output_handlers
else:
print("Not My Event")

def print_event(event):
print ("Normal Event", event.results)


def add_one(value):
return value + 1


with mp_event_loop.EventLoop(output_handlers=[print_my_event, print_event]) 作为循环:
loop.add_event(target=add_one, args=(1,) )
loop.add_event(target=add_one, args=(2,))
loop.add_event(MyEvent(3))
loop.add_event(target=add_one, args=(4,))
loop.add_event(MyEvent(5))


#不是我的事件
# 正常事件 2
# Not My Event
# Normal Event 3
# My Event [0, 1, 2]
# Not My Event
# Normal Event 5
# My Event [0, 1, 2, 3, 4]
```

## Pickling
如果酸洗很烦人然后你可以使用不同的多处理库。

EventLoop 使用 4 个类变量来创建正确的 Process 和 Thread 对象

* EventLoop.alive_event_class = multiprocessing.Event
* EventLoop.queue_class = multiprocessing.JoinableQueue
* EventLoop.event_loop_class = multiprocessing.Process
* EventLoop.consumer_loop_class = threading.Thread

`use`已提供功能以使此过程更容易。

```python
导入 mp_event_loop

导入线程
import multiprocess as mp

mp_event_loop.use(mp) # 这不会改变 consumer_loop_class
# mp_event_loop.use('multiprocess') # 也适用于字符串参数

# 或者

# 下面和 'use'
一样 mp_event_loop.EventLoop.alive_event_class = mp.Event
mp_event_loop.EventLoop.queue_class = mp.JoinableQueue
mp_event_loop.EventLoop.event_loop_class = mp.Process
mp_event_loop.EventLoop.consumer_loop_class = threading.Thread
```

### 酸洗问题
我的目标是扩展这个库以使用 async/等待。不幸的是,协程和生成器不能被
腌制。

最初,异步/等待事件循环失败,因为无法腌制生成器。我想创建另一个事件
循环,其中带有 yield 语句的函数将允许其他事件运行。虽然生成器不能被腌制,
但可以使用可以腌制的 _\_iter_\_ 和 _\_next_\_ 方法创建一个类。我创建了一个事件
循环(iter_event_loop.IterEventLoop),它将收集迭代器并交错迭代器。在 _\_next_\_ 被调用后
,一个不同的迭代器将执行添加更多并发。这只会使运行时间过长的迭代器不会占用
所有处理时间,并让其他迭代器事件在迭代之间运行。

我能够让 async/await 在多处理中工作。我使用了与缓存相同的技术,即
在类字典中注册异步函数并腌制注册的名称。在 unpickling 时,注册名称
用于检索异步函数并在其他进程中运行它。

项目详情


下载文件

下载适用于您平台的文件。如果您不确定要选择哪个,请了解有关安装包的更多信息。

源分布

mp_event_loop-1.5.2.tar.gz (24.4 kB 查看哈希

已上传 source