【问题标题】:Asynchronous message processing with async/await, RX and LINQ使用 async/await、RX 和 LINQ 进行异步消息处理
【发布时间】:2017-03-28 17:48:41
【问题描述】:

我将 Reactive Extensions 与 async/await 结合使用来简化我的套接字协议实现。当特定消息到达时必须执行一些操作(例如,向每个“ping”消息发送“pong”),还有一些方法我们必须异步等待一些特定的响应。以下示例说明了这一点:

private Subject<string> MessageReceived = new Subject<string>();

//this method gets called every time a message is received from socket
internal void OnReceiveMessage(string message)
{
    MessageReceived.OnNext(message);
    ProcessMessage(message);
}

public async Task<string> TestMethod()
{
    var expectedMessage = MessageReceived.Where(x => x.EndsWith("D") && x.EndsWith("F")).FirstOrDefaultAsync();
    await SendMessage("ABC");

    //some code...

    //if response we are waiting for comes before next row, we miss it
    return await expectedMessage;
}

TestMethod() 将“ABC”发送到套接字并在接收到“DEF”时继续(在此之前可能还有一些其他消息)。

这几乎可以工作,但有一个竞争条件。似乎这段代码在return await expectedMessage; 之前不会监听消息,这是一个问题,因为有时消息会在此之前到达。

【问题讨论】:

  • 你的意思是pong在等待expectedMessage的时候没有发送?在这种情况下,您应该在OnReceiveMessage 中处理ping 消息,而不是使用MessageReceived.OnNext() 存储它
  • 实际上我在示例中遗漏了那部分。它工作正常。在 ProcessMessages() 方法中,我发送 pong 以响应 ping 消息,但问题是在异步方法中等待特定事件(从套接字接收消息)

标签: c# .net linq async-await system.reactive


【解决方案1】:

FirstOrDefaultAsync 在这里不能很好地工作:直到 await 行它才会订阅,这会给你留下一个竞争条件(正如你指出的那样)。替换它的方法如下:

    var expectedMessage = MessageReceived
        .Where(x => x.EndsWith("D") && x.EndsWith("F"))
        .Take(1)
        .Replay(1)
        .RefCount();

    using (var dummySubscription = expectedMessage.Subscribe(i => {}))
    {
        await SendMessage("ABC");

        //Some code... goes here.

        return await expectedMessage;
    }

.Replay(1) 确保新订阅获取最新条目(假设存在)。它仅在有订阅者收听时才有效,因此dummySubscription

【讨论】:

    猜你喜欢
    • 2013-04-11
    • 2016-12-04
    • 1970-01-01
    • 1970-01-01
    • 2015-03-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多