Linq

Example
static async Task Main(string[] args)
{
    IObservable<int> observable = Observable.Create<int>(obs =>
    {
        return Task.Run(() =>
        {
            for (int i = 0; i < 10; i++)
            {
                obs.OnNext(i);
                Thread.Sleep(250);
            }
        });
    });

    observable.Subscribe(val => Console.WriteLine(val));

    await new TaskCompletionSource<object>().Task;
}

IObservables & Reactive LINQ Cheatsheet (Rx.NET)


Creating Observables

Observable.Return(value)             // Emits a single value
Observable.Range(start, count)       // Emits a sequence of numbers
Observable.Interval(TimeSpan)        // Emits numbers at intervals
Observable.Create<T>(observer => { /* ... */ })

Subscribing

IDisposable subscription = observable.Subscribe(
    onNext,     // Action<T>
    onError,    // Action<Exception>
    onComplete  // Action
);

Common Operators

.Select(x => x * 2)          // Map values
.Where(x => x > 5)           // Filter values
.Buffer(count)               // Group into lists of 'count'
.Buffer(TimeSpan, count)     // Buffer by time or count
.Take(n)                     // First n items
.TakeUntil(otherObservable)  // Until another observable emits
.Distinct()                  // Remove duplicates
.Merge(otherObservable)      // Combine streams
.CombineLatest(other)        // Latest values from both
.Zip(other, (x, y) => ...)   // Pairwise combine
.Delay(TimeSpan)             // Delay each value

Subjects (manual triggers)

var subject = new Subject<T>();
subject.OnNext(value);           // Emit value
subject.OnCompleted();           // Signal completion
subject.OnError(exception);      // Signal error

Schedulers

.ObserveOn(Scheduler)        // Observer runs on scheduler
.SubscribeOn(Scheduler)      // Source runs on scheduler

Disposing

subscription.Dispose();      // Stop receiving events

Error Handling

.Catch(otherObservable)
.Retry()

Creating Timers/Intervals

Observable.Timer(TimeSpan)          // One-off
Observable.Interval(TimeSpan)       // Repeated

Conversion

.ToEnumerable()    // To IEnumerable
.ToTask()          // To Task (async)

Useful Patterns


Docs: