Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
greenpool.py258 linesDownload Raw Back to eventlet
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 
codekingpro/portable-devtools · Team Ai