Diff of /buffering.py [000000] .. [70b6b3]

Switch to unified view

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