Have you heard about IObservable/IObserver support in Microsoft StreamInsight 1.1? Then you probably want to try it out. If this is your first incursion into the IObservable/IObserver pattern, this blog post is for you!
StreamInsight 1.1 introduced the ability to use IEnumerable and IObservable objects as event sources and sinks. The IEnumerable case is pretty straightforward, since many data collections are already surfacing as this type. This was already covered by Colin in his blog. Creating your own IObservable event source is a little more involved but no less exciting – here is a primer:
First, let’s look at a very simple Observable data source. All it does is publish an integer in regular time periods to its registered observers. (For more information on IObservable, see http://msdn.microsoft.com/en-us/library/dd990377.aspx ).
sealed class RandomSubject : IObservable<int>, IDisposable
{
private bool _done;
private readonly List<IObserver<int>> _observers;
private readonly Random _random;
private readonly object _sync;
private readonly Timer _timer;
private readonly int _timerPeriod;
/// <summary>
/// Random observable subject. It produces an integer in regular time periods.
/// </summary>
/// <param name="timerPeriod">Timer period (in milliseconds)</param>
public RandomSubject(int timerPeriod)
{
_done = false;
_observers = new List<IObserver<int>>();
_random = new Random();
_sync = new object();
_timer = new Timer(EmitRandomValue);
_timerPeriod = timerPeriod;
Schedule();
}
public IDisposable Subscribe(IObserver<int> observer)
{
lock (_sync)
{
_observers.Add(observer);
}
return new Subscription(this, observer);
}
public void OnNext(int value)
{
lock (_sync)
{
if (!_done)
{
foreach (var observer in _observers)
{
observer.OnNext(value);
}
}
}
}
public void OnError(Exception e)
{
lock (_sync)
{
foreach (var observer in _observers)
{
observer.OnError(e);
}
_done = true;
}
}
public void OnCompleted()
{
lock (_sync)
{
foreach (var observer in _observers)
{
observer.OnCompleted();
}
_done = true;
}
}
void IDisposable.Dispose()
{
_timer.Dispose();
}
private void Schedule()
{
lock (_sync)
{
if (!_done)
{
_timer.Change(_timerPeriod, Timeout.Infinite);
}
}
}
private void EmitRandomValue(object _)
{
var value = (int)(_random.NextDouble() * 100);
Console.WriteLine("[Observable]\t" + value);
OnNext(value);
Schedule();
}
private sealed class Subscription : IDisposable
{
private readonly RandomSubject _subject;
private IObserver<int> _observer;
public Subscription(RandomSubject subject, IObserver<int> observer)
{
_subject = subject;
_observer = observer;
}
public void Dispose()
{
IObserver<int> observer = _observer;
if (null != observer)
{
lock (_subject._sync)
{
_subject._observers.Remove(observer);
}
_observer = null;
}
}
}
}
So far, so good. Now let’s write a program that consumes data emitted by the observable as a stream of point events in a Streaminsight query. First, let’s define our payload type:
class Payload
{
public int Value { get; set; }
public override string ToString()
{
return "[StreamInsight]\tValue: " + Value.ToString();
}
}
Now, let’s write the program. First, we will instantiate the observable subject. Then we’ll use the ToPointStream() method to consume it as a stream. We can now write any query over the source - here, a simple pass-through query.
class Program
{
static void Main(string[] args)
{
Console.WriteLine("Starting observable source...");
using (var source = new RandomSubject(500))
{
Console.WriteLine("Started observable source.");
using (var server = Server.Create("Default"))
{
var application = server.CreateApplication("My Application");
var stream = source.ToPointStream(application,
e => PointEvent.CreateInsert(DateTime.Now, new Payload { Value = e }),
AdvanceTimeSettings.StrictlyIncreasingStartTime,
"Observable Stream");
var query = from e in stream
select e;
[...]
We’re done with consuming input and querying it! But you probably want to see the output of the query. Did you know you can turn a query into an observable subject as well? Let’s do precisely that, and exploit the Reactive Extensions for .NET (http://msdn.microsoft.com/en-us/devlabs/ee794896.aspx) to quickly visualize the output. Notice we’re subscribing “Console.WriteLine()” to the query, a pattern you may find useful for quick debugging of your queries. Reminder: you’ll need to install the Reactive Extensions for .NET (Rx for .NET Framework 4.0), and reference System.CoreEx and System.Reactive in your project.
[...]
Console.ReadLine();
Console.WriteLine("Starting query...");
using (query.ToObservable().Subscribe(Console.WriteLine))
{
Console.WriteLine("Started query.");
Console.ReadLine();
Console.WriteLine("Stopping query...");
}
Console.WriteLine("Stopped query.");
}
Console.ReadLine();
Console.WriteLine("Stopping observable source...");
source.OnCompleted();
}
Console.WriteLine("Stopped observable source.");
}
}
We hope this blog post gets you started. And for bonus points, you can go ahead and rewrite the observable source (the RandomSubject class) using the Reactive Extensions for .NET! The entire sample project is attached to this article.
Happy querying!
Regards, The StreamInsight Team