Rx: получение коллекции событий из "окна" - не выходит красиво избежать deadlock'а

у меня есть источник собитий source и он выдает по несколько событий за короткий промежуток времени, а потом молчит. я чуть не написал свою реализацию Window, нашел реализацию в доках, но наткнулся на deadlock.


вот как пытаюсь преобразовать IObservable<T> в IObservable<IList<T> (все события, которые получил в течении этого окна закинуть в список):

source.Window(TimeSpan.FromSeconds(1)).Select(obs => obs.ToListObservable())

по задумке, оно должно было блокировать вывод резултирующего события до окончания окна, но вместо этого просто делает deadlock.


подскажите как это реализовать без ручной обработки колбеков внутреннего IObservable


UPD:
хочу слоить 2 события в одно окно. чтобы первое событие начинало отсчет окна, а потом окно само закрывалось через TimeSpan. тестирую кодом:

var source = Observable
    .Interval(TimeSpan.FromMilliseconds(300));

var output = source.Window(TimeSpan.FromMilliseconds(700)).SelectMany(window => window.ToList());

using var subscription = output
    .Subscribe(list => Console.WriteLine(">>" + string.Join("; ", list)));
Console.ReadLine();

UPD2:
все еще пытаюсь настроить "Окно", чтобы работало как "выключатель": включаешь - начинает запись, выключаешь - заканчивает запись и появляется возможность включить.

var state = false;
var output = source.Delay(TimeSpan.FromMilliseconds(50))
    .Window(source.Where(next => state == false).Select(next=>
    {
        state = true;
        return DateTime.UtcNow + TimeSpan.FromMilliseconds(610);
    }), time => source.Where(next =>
    {
        var shellEnd = time < DateTime.UtcNow;
        if (shellEnd)
        {
            state = false;
        }
        return shellEnd;
    }))
    .SelectMany(window => window.ToList());

тут все еще есть проблемы:

  • иногда приходит ТРИ события в список вместо ДВУХ
  • задержка Delay часто не помогает. она должна немного откладывать вызов событий, чтобы окно открывалось/закрывалось прямо перед распространением события

Ответы (1 шт):

Автор решения: VladD

Функция Window очень мощная, позволяет разбить последовательность на мелкие кусочки («окна»), и контролировать вручную границы этих окон. Нам пригодится перегрузка с Func<IObservable<TWindowClosing>> windowClosingSelector (описание смотрите, например, тут, прокрутите до «RxNET Window» и распахните спойлер).

Эта перегрузка Window работает так. Первое окно открывается немедленно, и вызывается генератор закрывающего события: генератор возвращает IObservable<T>, и по приходу первого элемента окно закрывается. Тут же открывается новое окно, и вновь запрашивается генератор сгенерировать новое закрывающее событие. И так далее.

В качестве генератора закрывающих событий нам удобно использовать оператор Delay, применённый к исходной последовательности: он свдигает все элементы, начиная от текущего, на фиксированный временной промежуток. В результате при вызове генератора первым сгенерированным элементом будет следующий ещё не выданный на текущий момент времени элемент исходной последовательности, сдвинутый по времени на желаемую ширину окна. То есть окно закроется через заданный промежуток времени после прихода первого элемента.

Получив нужные окна, легко схлопнуть их в список про помощи простого .ToList().

Результирующий код:

var windowDuration = TimeSpan.FromMilliseconds(1500);
var output = source.Window(() => source.Delay(windowDuration))
                   .SelectMany(window => window.ToList());

Обратите внимание, что мы не можем закешировать source.Delay(windowDuration), т. к. нам нужно состояние именно на момент вызова.


Если ваш source холодный, и производит разные последовательности для разных подписчиков, нужно его превратить в горячий при помощи .Publish().RefCount(). Код основывается на том, что потоки в Window и в Delay совпадают.

→ Ссылка