paulb@146 | 1 | #!/usr/bin/env python |
paulb@146 | 2 | |
paulb@146 | 3 | """ |
paulb@146 | 4 | A simple example of parallel computation using persistent queues and |
paulb@146 | 5 | communications. |
paulb@146 | 6 | """ |
paulb@146 | 7 | |
paulb@146 | 8 | import pprocess |
paulb@146 | 9 | import time |
paulb@146 | 10 | #import random |
paulb@146 | 11 | import sys |
paulb@146 | 12 | |
paulb@146 | 13 | # Array size and a limit on the number of processes. |
paulb@146 | 14 | |
paulb@146 | 15 | N = 10 |
paulb@146 | 16 | limit = 2 # since N background processes will be used, this is reduced |
paulb@146 | 17 | delay = 1 |
paulb@146 | 18 | |
paulb@146 | 19 | # Work function and monitoring class. |
paulb@146 | 20 | |
paulb@146 | 21 | def calculate(i, j): |
paulb@146 | 22 | |
paulb@146 | 23 | """ |
paulb@146 | 24 | A supposedly time-consuming calculation on 'i' and 'j'. |
paulb@146 | 25 | """ |
paulb@146 | 26 | |
paulb@146 | 27 | #time.sleep(delay * random.random()) |
paulb@146 | 28 | time.sleep(delay) |
paulb@146 | 29 | return (i, j, i * N + j) |
paulb@146 | 30 | |
paulb@146 | 31 | # Background computation. |
paulb@146 | 32 | |
paulb@146 | 33 | def task(i): |
paulb@146 | 34 | |
paulb@146 | 35 | # Initialise the communications queue with a limit on the number of |
paulb@146 | 36 | # channels/processes. |
paulb@146 | 37 | |
paulb@146 | 38 | queue = pprocess.Queue(limit=limit) |
paulb@146 | 39 | |
paulb@146 | 40 | # Initialise an array. |
paulb@146 | 41 | |
paulb@146 | 42 | results = [0] * N |
paulb@146 | 43 | |
paulb@146 | 44 | # Wrap the calculate function and manage it. |
paulb@146 | 45 | |
paulb@146 | 46 | calc = queue.manage(pprocess.MakeParallel(calculate)) |
paulb@146 | 47 | |
paulb@146 | 48 | # Perform the work. |
paulb@146 | 49 | |
paulb@146 | 50 | print "Calculating..." |
paulb@146 | 51 | for j in range(0, N): |
paulb@146 | 52 | calc(i, j) |
paulb@146 | 53 | |
paulb@146 | 54 | # Store the results as they arrive. |
paulb@146 | 55 | |
paulb@146 | 56 | print "Finishing..." |
paulb@146 | 57 | for i, j, result in queue: |
paulb@146 | 58 | results[j] = result |
paulb@146 | 59 | |
paulb@146 | 60 | return i, results |
paulb@146 | 61 | |
paulb@146 | 62 | # Main program. |
paulb@146 | 63 | |
paulb@146 | 64 | if __name__ == "__main__": |
paulb@146 | 65 | |
paulb@146 | 66 | t = time.time() |
paulb@146 | 67 | |
paulb@146 | 68 | if "--reconnect" not in sys.argv: |
paulb@146 | 69 | |
paulb@146 | 70 | # Wrap the computation and manage it. |
paulb@146 | 71 | |
paulb@146 | 72 | ptask = pprocess.MakeParallel(task) |
paulb@146 | 73 | |
paulb@146 | 74 | for i in range(0, N): |
paulb@146 | 75 | |
paulb@146 | 76 | # Make a distinct callable for each part of the computation. |
paulb@146 | 77 | |
paulb@146 | 78 | ptask_i = pprocess.BackgroundCallable("task-%d.socket" % i, ptask) |
paulb@146 | 79 | |
paulb@146 | 80 | # Perform the work. |
paulb@146 | 81 | |
paulb@146 | 82 | ptask_i(i) |
paulb@146 | 83 | |
paulb@146 | 84 | # Discard the callable. |
paulb@146 | 85 | |
paulb@146 | 86 | del ptask |
paulb@146 | 87 | print "Discarded the callable." |
paulb@146 | 88 | |
paulb@146 | 89 | if "--start" not in sys.argv: |
paulb@146 | 90 | |
paulb@146 | 91 | # Open a queue and reconnect to the task. |
paulb@146 | 92 | |
paulb@146 | 93 | print "Opening a queue." |
paulb@146 | 94 | queue = pprocess.PersistentQueue() |
paulb@146 | 95 | for i in range(0, N): |
paulb@146 | 96 | queue.connect("task-%d.socket" % i) |
paulb@146 | 97 | |
paulb@146 | 98 | # Initialise an array. |
paulb@146 | 99 | |
paulb@146 | 100 | results = [0] * N |
paulb@146 | 101 | |
paulb@146 | 102 | # Wait for the results. |
paulb@146 | 103 | |
paulb@146 | 104 | print "Waiting for persistent results" |
paulb@146 | 105 | for i, result in queue: |
paulb@146 | 106 | results[i] = result |
paulb@146 | 107 | |
paulb@146 | 108 | # Show the results. |
paulb@146 | 109 | |
paulb@146 | 110 | for i in range(0, N): |
paulb@146 | 111 | for result in results[i]: |
paulb@146 | 112 | print result, |
paulb@146 | 113 | print |
paulb@146 | 114 | |
paulb@146 | 115 | print "Time taken:", time.time() - t |
paulb@146 | 116 | |
paulb@146 | 117 | # vim: tabstop=4 expandtab shiftwidth=4 |