|
a |
|
b/buffering.py |
|
|
1 |
import multiprocessing as mp |
|
|
2 |
import Queue |
|
|
3 |
import threading |
|
|
4 |
|
|
|
5 |
|
|
|
6 |
def buffered_gen_mp(source_gen, buffer_size=2): |
|
|
7 |
""" |
|
|
8 |
Generator that runs a slow source generator in a separate process. |
|
|
9 |
buffer_size: the maximal number of items to pre-generate (length of the buffer) |
|
|
10 |
""" |
|
|
11 |
if buffer_size < 2: |
|
|
12 |
raise RuntimeError("Minimal buffer size is 2!") |
|
|
13 |
|
|
|
14 |
buffer = mp.Queue(maxsize=buffer_size - 1) |
|
|
15 |
|
|
|
16 |
# the effective buffer size is one less, because the generation process |
|
|
17 |
# will generate one extra element and block until there is room in the buffer. |
|
|
18 |
|
|
|
19 |
def _buffered_generation_process(source_gen, buffer): |
|
|
20 |
for data in source_gen: |
|
|
21 |
buffer.put(data, block=True) |
|
|
22 |
buffer.put(None) # sentinel: signal the end of the iterator |
|
|
23 |
buffer.close() # unfortunately this does not suffice as a signal: if buffer.get() |
|
|
24 |
# was called and subsequently the buffer is closed, it will block forever. |
|
|
25 |
|
|
|
26 |
process = mp.Process(target=_buffered_generation_process, args=(source_gen, buffer)) |
|
|
27 |
process.start() |
|
|
28 |
|
|
|
29 |
for data in iter(buffer.get, None): |
|
|
30 |
yield data |
|
|
31 |
|
|
|
32 |
|
|
|
33 |
def buffered_gen_threaded(source_gen, buffer_size=5): |
|
|
34 |
""" |
|
|
35 |
Generator that runs a slow source generator in a separate thread. Beware of the GIL! |
|
|
36 |
buffer_size: the maximal number of items to pre-generate (length of the buffer) |
|
|
37 |
""" |
|
|
38 |
if buffer_size < 2: |
|
|
39 |
raise RuntimeError("Minimal buffer size is 2!") |
|
|
40 |
|
|
|
41 |
buffer = Queue.Queue(maxsize=buffer_size - 1) |
|
|
42 |
|
|
|
43 |
# the effective buffer size is one less, because the generation process |
|
|
44 |
# will generate one extra element and block until there is room in the buffer. |
|
|
45 |
|
|
|
46 |
def _buffered_generation_thread(source_gen, buffer): |
|
|
47 |
for data in source_gen: |
|
|
48 |
buffer.put(data, block=True) |
|
|
49 |
buffer.put(None) # sentinel: signal the end of the iterator |
|
|
50 |
|
|
|
51 |
thread = threading.Thread(target=_buffered_generation_thread, args=(source_gen, buffer)) |
|
|
52 |
thread.daemon = True |
|
|
53 |
thread.start() |
|
|
54 |
|
|
|
55 |
for data in iter(buffer.get, None): |
|
|
56 |
yield data |