BehaviorSubject per group with GroupBy and Switch()

Viewed 168

I have code which would need to have GroupBy and would need a unique BehaviorSubject per group of Switch().

We have a stream of stock market values that we group by Symbol and perform level crossing across a number of levels (defined by a BehaviorSubject and a switch to always use the latest values).

So I need to go from this:

var feed = new Subject<double>();
var levels = new BehaviorSubject<double[]>(new[] { 400.0, 500.0, 600.0, 700.0 });

levels
    .Select(thresholds => feed
        .Buffer(2, 1)
        .Where(x => x.Count == 2)
        .Select(x => new { LevelsCrossed = thresholds.GetCrossovers(x[0], x[1]), Previous = x[0], Current = x[1] })
        .Where(x => x.LevelsCrossed.Any())
        .SelectMany(x => x.LevelsCrossed.Select(level => new ThresholdCrossedEvent(level, x.Previous, x.Current))))
    .Switch()
    .Subscribe(x => Console.WriteLine(JsonConvert.SerializeObject(x)));

And adapt the above to take a stream of Tick below and group by Symbol, each with its own level threshold detection on each grouped Value.

class Tick
{
    public string Symbol { get; set; } // The name.
    public decimal Value { get; set; } // The value.
}

Outline:

  1. Take Market data
  2. Group by Symbol
  3. Alert on levels (depending on group name, using a dictionary of BehaviorSubject)
  4. Output
  5. Use Switch() to always use latest values from the dictionary

With a naive implementation I have a wrapper class (ReactiveSymbolFeed below), however blurring non-reactive and reactive code can introduce potential concurrency issues that reactive extensions otherwise deals neatly with.

Questions please:

Am I introducing any side effects, or will this cause issue at scale (say 100,000 messages per second across 2,000 groups)?

Since we have many groups each with their own BehaviorSubject that needs Switch() - can we rewrite our Reactive Extensions statement block to include the thresholds levels per symbol group, or is the above wrapper class the right way to do this?

Further context and the wrapper class solution

Instead I create a ReactiveSymbolFeed wrapper that will form the value part of a dictionary per symbol key.

class ReactiveSymbolFeed
{
    readonly BehaviorSubject<double[]> levels;
    readonly Subject<double> feed;

    public ReactiveSymbolFeed(double[] levels)
    {
        this.feed = new Subject<double>();
        this.levels = new BehaviorSubject<double[]>(levels);

        this.levels
            .Select(thresholds => this.feed
            .Buffer(2, 1)
            .Where(x => x.Count == 2)
            .Select(x => new { LevelsCrossed = thresholds.GetCrossovers(x[0], x[1]), Previous = x[0], Current = x[1] })
            .Where(x => x.LevelsCrossed.Any())
            .SelectMany(x => x.LevelsCrossed.Select(level => new ThresholdCrossedEvent(level, x.Previous, x.Current))))
         .Switch()
         .DistinctUntilChanged(x => x.Threshold)
         .Subscribe(x => Console.WriteLine(JsonConvert.SerializeObject(x)));
    }

    public void OnNext(double value) => this.feed.OnNext(value);

    public void UpdateThresholds(double[] levels) => this.levels.OnNext(levels);
}

And then use with the below:

// Setup the detection thresholds per Symbol - each Symbol has 1 set of thresholds
var dictionary = new Dictionary<string, ReactiveSymbolFeed>();            
dictionary.Add("AAPL", new ReactiveSymbolFeed(new[] { 120.0, 125.0, 130.0 }));
dictionary.Add("VXX", new ReactiveSymbolFeed(new[] { 10.5, 15, 18.5, 20 }));

// Create some test tick data.
var ticks = new[]
{
    new Tick { Symbol = "AAPL", Value = 119.0 },
    new Tick { Symbol = "VXX", Value = 10.3 },
    new Tick { Symbol = "VXX", Value = 10.8 },
    new Tick { Symbol = "AAPL", Value = 121.0 },
    new Tick { Symbol = "AAPL", Value = 121.0 }
    // Followed by many other differnet Symbols and Values
};

// Loop through test data and dispatch it.
foreach(var tick in ticks)
{
    if(dictionary.TryGetValue(tick.Symbol, out var value))
        value.OnNext(tick.Value);
}
0 Answers
Related