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 valueSubjects (manual triggers)
var subject = new Subject<T>();
subject.OnNext(value); // Emit value
subject.OnCompleted(); // Signal completion
subject.OnError(exception); // Signal errorSchedulers
.ObserveOn(Scheduler) // Observer runs on scheduler
.SubscribeOn(Scheduler) // Source runs on schedulerDisposing
subscription.Dispose(); // Stop receiving eventsError Handling
.Catch(otherObservable)
.Retry()Creating Timers/Intervals
Observable.Timer(TimeSpan) // One-off
Observable.Interval(TimeSpan) // RepeatedConversion
.ToEnumerable() // To IEnumerable
.ToTask() // To Task (async)Useful Patterns
- Graceful shutdown: Use .TakeUntil(stopSignal) to complete stream on signal.
- Batching: .Buffer(TimeSpan, count) to emit batches.
- Threading: Use .ObserveOn(ThreadPoolScheduler.Instance) for background processing.
Docs: