codekingpro/portable-devtools
114k
1import traceback2 3import eventlet4from eventlet import queue5from eventlet.support import greenlets as greenlet6import six7 8__all__ = ['GreenPool', 'GreenPile']9 10DEBUG = True11 12 13class GreenPool(object):14 """The GreenPool class is a pool of green threads.15 """16 17 def __init__(self, size=1000):18 try:19 size = int(size)20 except ValueError as e:21 msg = 'GreenPool() expect size :: int, actual: {0} {1}'.format(type(size), str(e))22 raise TypeError(msg)23 if size < 0:24 msg = 'GreenPool() expect size >= 0, actual: {0}'.format(repr(size))25 raise ValueError(msg)26 self.size = size27 self.coroutines_running = set()28 self.sem = eventlet.Semaphore(size)29 self.no_coros_running = eventlet.Event()30 31 def resize(self, new_size):32 """ Change the max number of greenthreads doing work at any given time.33 34 If resize is called when there are more than *new_size* greenthreads35 already working on tasks, they will be allowed to complete but no new36 tasks will be allowed to get launched until enough greenthreads finish37 their tasks to drop the overall quantity below *new_size*. Until38 then, the return value of free() will be negative.39 """40 size_delta = new_size - self.size41 self.sem.counter += size_delta42 self.size = new_size43 44 def running(self):45 """ Returns the number of greenthreads that are currently executing46 functions in the GreenPool."""47 return len(self.coroutines_running)48 49 def free(self):50 """ Returns the number of greenthreads available for use.51 52 If zero or less, the next call to :meth:`spawn` or :meth:`spawn_n` will53 block the calling greenthread until a slot becomes available."""54 return self.sem.counter55 56 def spawn(self, function, *args, **kwargs):57 """Run the *function* with its arguments in its own green thread.58 Returns the :class:`GreenThread <eventlet.GreenThread>`59 object that is running the function, which can be used to retrieve the60 results.61 62 If the pool is currently at capacity, ``spawn`` will block until one of63 the running greenthreads completes its task and frees up a slot.64 65 This function is reentrant; *function* can call ``spawn`` on the same66 pool without risk of deadlocking the whole thing.67 """68 # if reentering an empty pool, don't try to wait on a coroutine freeing69 # itself -- instead, just execute in the current coroutine70 current = eventlet.getcurrent()71 if self.sem.locked() and current in self.coroutines_running:72 # a bit hacky to use the GT without switching to it73 gt = eventlet.greenthread.GreenThread(current)74 gt.main(function, args, kwargs)75 return gt76 else:77 self.sem.acquire()78 gt = eventlet.spawn(function, *args, **kwargs)79 if not self.coroutines_running:80 self.no_coros_running = eventlet.Event()81 self.coroutines_running.add(gt)82 gt.link(self._spawn_done)83 return gt84 85 def _spawn_n_impl(self, func, args, kwargs, coro):86 try:87 try:88 func(*args, **kwargs)89 except (KeyboardInterrupt, SystemExit, greenlet.GreenletExit):90 raise91 except:92 if DEBUG:93 traceback.print_exc()94 finally:95 if coro is None:96 return97 else:98 coro = eventlet.getcurrent()99 self._spawn_done(coro)100 101 def spawn_n(self, function, *args, **kwargs):102 """Create a greenthread to run the *function*, the same as103 :meth:`spawn`. The difference is that :meth:`spawn_n` returns104 None; the results of *function* are not retrievable.105 """106 # if reentering an empty pool, don't try to wait on a coroutine freeing107 # itself -- instead, just execute in the current coroutine108 current = eventlet.getcurrent()109 if self.sem.locked() and current in self.coroutines_running:110 self._spawn_n_impl(function, args, kwargs, None)111 else:112 self.sem.acquire()113 g = eventlet.spawn_n(114 self._spawn_n_impl,115 function, args, kwargs, True)116 if not self.coroutines_running:117 self.no_coros_running = eventlet.Event()118 self.coroutines_running.add(g)119 120 def waitall(self):121 """Waits until all greenthreads in the pool are finished working."""122 assert eventlet.getcurrent() not in self.coroutines_running, \123 "Calling waitall() from within one of the " \124 "GreenPool's greenthreads will never terminate."125 if self.running():126 self.no_coros_running.wait()127 128 def _spawn_done(self, coro):129 self.sem.release()130 if coro is not None:131 self.coroutines_running.remove(coro)132 # if done processing (no more work is waiting for processing),133 # we can finish off any waitall() calls that might be pending134 if self.sem.balance == self.size:135 self.no_coros_running.send(None)136 137 def waiting(self):138 """Return the number of greenthreads waiting to spawn.139 """140 if self.sem.balance < 0:141 return -self.sem.balance142 else:143 return 0144 145 def _do_map(self, func, it, gi):146 for args in it:147 gi.spawn(func, *args)148 gi.done_spawning()149 150 def starmap(self, function, iterable):151 """This is the same as :func:`itertools.starmap`, except that *func* is152 executed in a separate green thread for each item, with the concurrency153 limited by the pool's size. In operation, starmap consumes a constant154 amount of memory, proportional to the size of the pool, and is thus155 suited for iterating over extremely long input lists.156 """157 if function is None:158 function = lambda *a: a159 # We use a whole separate greenthread so its spawn() calls can block160 # without blocking OUR caller. On the other hand, we must assume that161 # our caller will immediately start trying to iterate over whatever we162 # return. If that were a GreenPile, our caller would always see an163 # empty sequence because the hub hasn't even entered _do_map() yet --164 # _do_map() hasn't had a chance to spawn a single greenthread on this165 # GreenPool! A GreenMap is safe to use with different producer and166 # consumer greenthreads, because it doesn't raise StopIteration until167 # the producer has explicitly called done_spawning().168 gi = GreenMap(self.size)169 eventlet.spawn_n(self._do_map, function, iterable, gi)170 return gi171 172 def imap(self, function, *iterables):173 """This is the same as :func:`itertools.imap`, and has the same174 concurrency and memory behavior as :meth:`starmap`.175 176 It's quite convenient for, e.g., farming out jobs from a file::177 178 def worker(line):179 return do_something(line)180 pool = GreenPool()181 for result in pool.imap(worker, open("filename", 'r')):182 print(result)183 """184 return self.starmap(function, six.moves.zip(*iterables))185 186 187class GreenPile(object):188 """GreenPile is an abstraction representing a bunch of I/O-related tasks.189 190 Construct a GreenPile with an existing GreenPool object. The GreenPile will191 then use that pool's concurrency as it processes its jobs. There can be192 many GreenPiles associated with a single GreenPool.193 194 A GreenPile can also be constructed standalone, not associated with any195 GreenPool. To do this, construct it with an integer size parameter instead196 of a GreenPool.197 198 It is not advisable to iterate over a GreenPile in a different greenthread199 than the one which is calling spawn. The iterator will exit early in that200 situation.201 """202 203 def __init__(self, size_or_pool=1000):204 if isinstance(size_or_pool, GreenPool):205 self.pool = size_or_pool206 else:207 self.pool = GreenPool(size_or_pool)208 self.waiters = queue.LightQueue()209 self.counter = 0210 211 def spawn(self, func, *args, **kw):212 """Runs *func* in its own green thread, with the result available by213 iterating over the GreenPile object."""214 self.counter += 1215 try:216 gt = self.pool.spawn(func, *args, **kw)217 self.waiters.put(gt)218 except:219 self.counter -= 1220 raise221 222 def __iter__(self):223 return self224 225 def next(self):226 """Wait for the next result, suspending the current greenthread until it227 is available. Raises StopIteration when there are no more results."""228 if self.counter == 0:229 raise StopIteration()230 return self._next()231 __next__ = next232 233 def _next(self):234 try:235 return self.waiters.get().wait()236 finally:237 self.counter -= 1238 239 240# this is identical to GreenPile but it blocks on spawn if the results241# aren't consumed, and it doesn't generate its own StopIteration exception,242# instead relying on the spawning process to send one in when it's done243class GreenMap(GreenPile):244 def __init__(self, size_or_pool):245 super(GreenMap, self).__init__(size_or_pool)246 self.waiters = queue.LightQueue(maxsize=self.pool.size)247 248 def done_spawning(self):249 self.spawn(lambda: StopIteration())250 251 def next(self):252 val = self._next()253 if isinstance(val, StopIteration):254 raise val255 else:256 return val257 __next__ = next258 