Skip to content

Filtering

A stream often carries more than you need. A filtering operator drops the values you do not want and passes the rest through unchanged.

Every operator here keeps values as they are. None of them changes a value. If you need to change values as well as drop them, use Choose, which is on the transformation page.

using ReactiveUI.Primitives;
using ReactiveUI.Primitives.Signals;

Your first filter

You have a stream of temperature readings. You only care about the ones above freezing, and only when the reading actually changes.

1. Start with the stream.

IObservable<int> readings = Signal.FromEnumerable([-5, 3, 3, 8, -1, 8]);

2. Drop the ones you do not want. Where passes through only the values your test accepts.

IObservable<int> aboveFreezing = readings.Where(x => x > 0);

After this step the stream carries 3, 3, 8, 8.

3. Drop repeats. Unique drops a value when it is the same as the one just before it.

IObservable<int> changes = aboveFreezing.Unique();

4. Subscribe.

changes.Subscribe(x => Console.Write($"{x} "));

Output:

3 8

Trace it one step at a time:

AfterStream carries
the source-5, 3, 3, 8, -1, 8
Where(x => x > 0)3, 3, 8, 8
Unique()3, 8

The two 3s collapsed because they sat next to each other. The two 8s collapsed too, and that is the part worth slowing down on. In the source they had -1 between them. Where dropped the -1 first, so by the time Unique saw the stream the two 8s were neighbours.

Order matters. Unique only ever sees what the operator before it passed on. Put the two calls the other way round and you get 3, 8, 8, because Unique would then see the -1 sitting between the 8s.

Testing each value

Where

Where passes through the values your test accepts.

Input: 1, 2, 3, 4, 5, 6

IObservable<int> evens = Signal.FromEnumerable([1, 2, 3, 4, 5, 6])
                               .Where(x => x % 2 == 0);

evens.Subscribe(x => Console.Write($"{x} "));

Output:

2 4 6

KeepWith

KeepWith passes a state value into the test as its first argument, so you can mark the lambda static and it allocates no closure. See mark lambdas static.

Input: 1, 2, 3, 4, 5, 6

var minimum = 4;

IObservable<int> large = Signal.FromEnumerable([1, 2, 3, 4, 5, 6])
                               .KeepWith(minimum, static (min, x) => x >= min);

large.Subscribe(x => Console.Write($"{x} "));

Output:

4 5 6

The result matches Where(x => x >= minimum). Only the allocation differs.

Dropping nulls and wrong types

KeepNotNull

KeepNotNull drops nulls. It also narrows the type, so what comes out is not nullable and you do not need a null check downstream.

Input: "Ada", null, "Grace"

IObservable<string?> maybeNames = Signal.FromEnumerable<string?>(["Ada", null, "Grace"]);

IObservable<string> names = maybeNames.KeepNotNull();

names.Subscribe(x => Console.WriteLine(x.Length));   // no null check needed

Output:

3
5

It works on reference types.

OfType

OfType keeps only the values that are of the type you name, and drops the rest.

Input: 1, "two", 3, "four" as object

IObservable<object?> mixed = Signal.FromEnumerable<object?>([1, "two", 3, "four"]);

IObservable<int> numbers = mixed.OfType<int>();

numbers.Subscribe(x => Console.Write($"{x} "));

Output:

1 3

OfType against Cast

Both narrow a stream of object? to one type. They disagree about a value that does not fit.

Input: 1, "two", 3

OperatorOutputOn a value of the wrong type
OfType<int>()1 3Drops it and carries on.
Cast<int>()1, then an errorEnds the stream with InvalidCastException.

Use OfType when mixed types are expected. Use Cast when a wrong type means a bug you want to hear about. Cast is on the transformation page.

Taking and skipping by count

Take

Take passes through at most the first few values, then completes. It does not wait for the source to end.

Input: 1, 2, 3, 4, 5

IObservable<int> firstThree = Signal.FromEnumerable([1, 2, 3, 4, 5]).Take(3);

firstThree.Subscribe(x => Console.Write($"{x} "), () => Console.WriteLine("| done"));

Output:

1 2 3 | done

Take on an endless stream gives you a stream that ends, which is how you turn a timer into a countdown.

Signal.Every(TimeSpan.FromSeconds(1)).Take(3).Subscribe(x => Console.WriteLine(x));

Output, one per second, then it stops on its own:

0
1
2

Skip

Skip ignores the first few values and passes through everything after them.

Input: 1, 2, 3, 4, 5

IObservable<int> rest = Signal.FromEnumerable([1, 2, 3, 4, 5]).Skip(2);

rest.Subscribe(x => Console.Write($"{x} "));

Output:

3 4 5

Taking and skipping by test

TakeWhile

TakeWhile passes values through while your test holds. The first value that fails the test ends the stream. That value is not sent.

Input: 1, 2, 3, 10, 4, 5

IObservable<int> small = Signal.FromEnumerable([1, 2, 3, 10, 4, 5])
                               .TakeWhile(x => x < 5);

small.Subscribe(x => Console.Write($"{x} "), () => Console.WriteLine("| done"));

Output:

1 2 3 | done

The 4 and 5 never arrive, even though they pass the test. TakeWhile stopped at the 10.

SkipWhile

SkipWhile drops values while your test holds, then passes through everything after that. Once it starts passing values, it never stops.

Input: 1, 2, 3, 10, 4, 5

IObservable<int> fromFirstLarge = Signal.FromEnumerable([1, 2, 3, 10, 4, 5])
                                        .SkipWhile(x => x < 5);

fromFirstLarge.Subscribe(x => Console.Write($"{x} "));

Output:

10 4 5

The 4 and 5 arrive here, because the test stopped being consulted once 10 failed it.

TakeWhile against Where

This is worth seeing on one input, because the two are easy to mix up.

Input: 1, 2, 3, 10, 4, 5, test x < 5

OperatorOutputRule
Where(x => x < 5)1 2 3 4 5Tests every value, forever.
TakeWhile(x => x < 5)1 2 3Stops the stream at the first failure.

Stopping on another event

TakeUntil

TakeUntil passes values through until a second stream sends anything. Then it completes.

Input: values every 100ms, a stop stream that fires at 250ms

IObservable<long> ticks = Signal.Every(TimeSpan.FromMilliseconds(100));
IObservable<long> stop = Signal.After(TimeSpan.FromMilliseconds(250));

ticks.TakeUntil(stop).Subscribe(x => Console.Write($"{x} "), () => Console.WriteLine("| done"));

Output:

0 1 | done

Only a value stops it. A stop stream that completes without sending anything does not stop it. Use Signal.After rather than Signal.Empty when you want a deadline.

A second form takes a CancellationToken instead of a stream.

IObservable<long> untilCancelled = ticks.TakeUntil(cancellationToken);

When the token is cancelled the stream completes normally. It does not error, so your completion callback runs and your error callback does not.

Dropping duplicates

Two operators drop duplicates and they mean different things by "duplicate".

Distinct

Distinct remembers every value it has ever sent, and drops any value it has seen before.

Input: 1, 1, 2, 1, 3, 2

IObservable<int> seenOnce = Signal.FromEnumerable([1, 1, 2, 1, 3, 2]).Distinct();

seenOnce.Subscribe(x => Console.Write($"{x} "));

Output:

1 2 3

Because it remembers every value, Distinct uses more memory the longer the stream runs. On an endless stream of many different values, that memory is never released.

Pass an IEqualityComparer<T> to decide what counts as equal.

IObservable<string> ignoringCase = names.Distinct(StringComparer.OrdinalIgnoreCase);

DistinctBy

DistinctBy does the same, comparing a key you pick out of each value rather than the whole value.

Input: three people, two of them in London

IObservable<Person> firstPerCity = people.DistinctBy(person => person.City);

firstPerCity.Subscribe(p => Console.WriteLine($"{p.Name} ({p.City})"));

Output:

Ada (London)
Grace (Paris)

Katherine, also in London, is dropped because London was already seen. It also takes a comparer for the key.

Unique

Unique compares each value only with the one immediately before it. It has no memory beyond that, so it costs nothing to run on an endless stream.

Input: 1, 1, 2, 1, 3, 2

IObservable<int> changes = Signal.FromEnumerable([1, 1, 2, 1, 3, 2]).Unique();

changes.Subscribe(x => Console.Write($"{x} "));

Output:

1 2 1 3 2

Only the second 1 was dropped, because it sat next to the first.

UniqueBy

UniqueBy compares a key from each value with the key from the value before.

Input: readings from the same sensor, then a different one, then the first again

IObservable<Reading> sensorChanges = readings.UniqueBy(r => r.SensorId);

sensorChanges.Subscribe(r => Console.WriteLine(r.SensorId));

Output:

A
B
A

It takes an optional comparer for the key.

Distinct against Unique

The same input through both, side by side.

Input: 1, 1, 2, 1, 3, 2

OperatorOutputRuleMemory
Distinct()1 2 3Drops a value seen anywhere earlier.Grows with the number of different values.
Unique()1 2 1 3 2Drops a value only when it repeats the one before.One value.

Use Unique to react to change, such as a property that keeps reporting the same value. Use Distinct to remove duplicates from a finite set, such as ids from a finished query. On an endless stream, Unique is almost always the one you want.

Dropping everything

IgnoreValues

IgnoreValues drops every value and passes on only the ending. Use it when you care that a job finished, not what it produced.

Input: 1, 2, 3, then completion

IObservable<int> finished = Signal.FromEnumerable([1, 2, 3]).IgnoreValues();

finished.Subscribe(
    x => Console.WriteLine($"value: {x}"),
    () => Console.WriteLine("done"));

Output:

done

An error still comes through, so this is a clean way to wait for a job and hear about failure.

Filling in an empty stream

DefaultIfEmpty

DefaultIfEmpty sends one value if the source completes without ever sending anything. A source that does send values is passed through untouched.

Input: nothing, then completion

IObservable<int> nothing = Signal.Empty<int>();

nothing.DefaultIfEmpty().Subscribe(x => Console.WriteLine(x));

Output:

0

Give it your own fallback instead of default:

Signal.Empty<string>()
      .DefaultIfEmpty("no results")
      .Subscribe(x => Console.WriteLine(x));

Output:

no results

With a source that does send something:

Input: 1, 2

Signal.FromEnumerable([1, 2]).DefaultIfEmpty(99).Subscribe(x => Console.Write($"{x} "));

Output:

1 2

The fallback never appeared, because the source was not empty.

Every filtering operator at a glance

OperatorSecond nameDrops
WhereKeepValues your test rejects.
KeepWithWhereWithThe same, with no closure allocated.
KeepNotNullWhereNotNullNulls, and narrows the type.
OfType<T>KeepType<T>Values that are not a T.
TakeEverything after the first n.
SkipThe first n.
TakeWhileEverything from the first failure onward.
SkipWhileValues before the first failure.
TakeUntilEverything after another stream fires.
DistinctValues seen anywhere earlier.
DistinctByValues whose key was seen earlier.
UniqueDistinctUntilChangedValues equal to the one before.
UniqueByDistinctUntilChangedByValues whose key equals the key before.
IgnoreValuesIgnoreElementsEvery value.
DefaultIfEmptyNothing. It adds a value to an empty stream.

The types behind these operators

KeepSignal, KeepWithSignal, KeepNotNullSignal, KeepTypeSignal, UniqueSignal and IgnoreValuesSignal are public classes in ReactiveUI.Primitives.Advanced, and each takes its source through the constructor.

using ReactiveUI.Primitives.Advanced;

// the operator
IObservable<int> a = source.Unique();

// the same thing, built directly
IObservable<int> b = new UniqueSignal<int>(source, EqualityComparer<int>.Default);

Calling the operator is the normal path. Construct the type when you are writing an operator of your own and want to place one inside it.

The two package flavours

Every operator here ships twice. ReactiveUI.Primitives puts them under ReactiveUI.Primitives.*. ReactiveUI.Primitives.Reactive is the same source compiled against System.Reactive, and puts them under ReactiveUI.Primitives.Reactive.*. They behave the same in both.