1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68
|
import asyncio
from asyncio import Future
import reactivex
from reactivex import Observable
from reactivex import operators as ops
from reactivex.scheduler.eventloop import AsyncIOScheduler
def to_async_iterable():
def _to_async_iterable(source: Observable):
class AIterable:
def __aiter__(self):
class AIterator:
def __init__(self):
self.notifications = []
self.future = Future()
source.pipe(ops.materialize()).subscribe(self.on_next)
def feeder(self):
if not self.notifications or self.future.done():
return
notification = self.notifications.pop(0)
dispatch = {
"N": lambda: self.future.set_result(notification.value),
"E": lambda: self.future.set_exception(
notification.exception
),
"C": lambda: self.future.set_exception(StopAsyncIteration),
}
dispatch[notification.kind]()
def on_next(self, notification):
self.notifications.append(notification)
self.feeder()
async def __anext__(self):
self.feeder()
value = await self.future
self.future = Future()
return value
return AIterator()
return AIterable()
return _to_async_iterable
async def go(loop):
scheduler = AsyncIOScheduler(loop)
ai = reactivex.range(0, 10, scheduler=scheduler).pipe(to_async_iterable())
async for x in ai:
print("got %s" % x)
def main():
loop = asyncio.get_event_loop()
loop.run_until_complete(go(loop))
if __name__ == "__main__":
main()
|