Closed austinnichols101 closed 3 years ago
I ended up creating a AsyncSubject() -> AsyncAnonymousObserver()
stream and then looping to "publish" the records:
for record in records()
await stream.asend(value=record
I would be nice to have a from_async_generator
operation :) but the approach above works just fine and is very readable.
After adding some logging, I was also able to see "Implicit synchronous back-pressure ™" in action. Very nice!
How can I create an
AsyncObservable
from anasync_generator
using aioreactive?results in an Exception:
TypeError("'async_generator' object is not iterable")