这是我的Interval定义:
m_interval = Observable.Interval(TimeSpan.FromSeconds(5), m_schedulerProvider.EventLoop)
.ObserveOn(m_schedulerProvider.EventLoop)
.Select(l => Observable.FromAsync(DoWork))
.Concat()
.Subscribe();
Run Code Online (Sandbox Code Playgroud)
在上面的代码中,我提供ISchedulerin Interval和ObserveOnfrom,SchedulerProvider以便我可以更快地进行单元测试(TestScheduler.AdvanceBy).另外,DoWork是一种async方法.
在我的特定情况下,我希望DoWork每5秒调用一次该函数.这里的问题是我希望5秒是DoWork另一个结束和开始之间的时间.因此,如果DoWork执行时间超过5秒,假设为10秒,则第一次呼叫为5秒,第二次呼叫为15秒.
不幸的是,以下测试证明它不像那样:
[Fact]
public void MultiPluginStatusHelperShouldWaitForNextQuery()
{
m_queryHelperMock
.Setup(x => x.CustomQueryAsync())
.Callback(() => Thread.Sleep(10000))
.Returns(Task.FromResult(new QueryCompletedEventData()))
.Verifiable()
;
var multiPluginStatusHelper = m_container.GetInstance<IMultiPluginStatusHelper>();
multiPluginStatusHelper.MillisecondsInterval = 5000;
m_testSchedulerProvider.EventLoopScheduler.AdvanceBy(TimeSpan.FromMilliseconds(5000).Ticks);
m_testSchedulerProvider.EventLoopScheduler.AdvanceBy(TimeSpan.FromMilliseconds(5000).Ticks);
m_queryHelperMock.Verify(x => x.CustomQueryAsync(), Times.Once);
}
Run Code Online (Sandbox Code Playgroud)
在DoWork调用CustomQueryAsync和测试失败说的就是被称为两次.它应该只被调用一次,因为强制延迟.Callback(() => Thread.Sleep(1000)).
我在这做错了什么?
我的实际实现来自这个例子.
我有一个asnyc函数,我想在IObservable序列中的每个观察中调用,一次限制一个事件的传递.消费者希望飞行中不超过一条消息; 如果我理解正确的话,这也是RX合约.
考虑这个样本:
static void Main() {
var ob = Observable.Interval(TimeSpan.FromMilliseconds(100));
//var d = ob.Subscribe(async x => await Consume(x)); // Does not rate-limit.
var d = ob.Subscribe(x => Consume(x).Wait());
Thread.Sleep(10000);
d.Dispose();
}
static async Task<Unit> Consume(long count) {
Console.WriteLine($"Consuming {count} on thread {Thread.CurrentThread.ManagedThreadId}");
await Task.Delay(750);
Console.WriteLine($"Returning on thread {Thread.CurrentThread.ManagedThreadId}");
return Unit.Default;
}
Run Code Online (Sandbox Code Playgroud)
该Consume函数伪造750毫秒的处理时间,并ob每100毫秒产生一次事件.上面的代码有效,但调用task.Wait()随机线程.如果我在注释掉的第3行中订阅,则以与生成事件Consume相同的速率调用ob(我甚至无法理解Subscribe我在此注释语句中使用的重载,因此它可能是无意义的).
那么如何从可观察序列到async函数一次正确地传递一个事件呢?
给定一个同步观察者,有没有办法做到这一点:
observable.SubscribeAsync(observer);
Run Code Online (Sandbox Code Playgroud)
并且observer异步调用所有方法或者是在创建观察者时必须处理的内容吗?
c# multithreading asynchronous reactive-programming system.reactive