【问题标题】:How to avoid this race condition in Reactive Extensions如何在响应式扩展中避免这种竞争条件
【发布时间】:2012-09-10 11:04:34
【问题描述】:

我在方法中有这样的代码:

ISubject<Message> messages = new ReplaySubject<Message>(messageTimeout);

public void HandleNext(string clientId, Action<object> callback)
{
    messages.Where(message => !message.IsHandledBy(clientId))
            .Take(1)
            .Subscribe(message =>
                       {
                           callback(message.Message);
                           message.MarkAsHandledBy(clientId);
                       });
}

编码它的 rx'y 方法是什么,这样MarkAsHandledBy()IsHandledBy() 之间就不会在多个并发调用 HandleNext() 时发生竞争?

编辑:

这是用于长轮询。 HandleNext() 为每个 Web 请求调用。请求只能处理一条消息,然后返回给客户端。下一个请求接收下一条消息,依此类推。

完整的代码(当然仍在进行中)是这样的:

public class Queue
{
    readonly ISubject<MessageWrapper> messages;

    public Queue() : this(TimeSpan.FromSeconds(30)) {}

    public Queue(TimeSpan messageTimeout)
    {
        messages = new ReplaySubject<MessageWrapper>(messageTimeout);
    }

    public void Send(string channel, object message)
    {
        messages.OnNext(new MessageWrapper(new List<string> {channel}, message));
    }

    public void ReceiveNext(string clientId, string channel, Action<object> callback)
    {
        messages
            .Where(message => message.Channels.Contains(channel) && !message.IsReceivedBy(clientId))
            .Take(1)
            .Subscribe(message =>
                       {
                           callback(message.Message);
                           message.MarkAsReceivedFor(clientId);
                       });
    }

    class MessageWrapper
    {
        readonly List<string> receivers;

        public MessageWrapper(List<string> channels, object message)
        {
            receivers = new List<string>();
            Channels = channels;
            Message = message;
        }

        public List<string> Channels { get; private set; }
        public object Message { get; private set; }

        public void MarkAsReceivedFor(string clientId)
        {
            receivers.Add(clientId);
        }

        public bool IsReceivedBy(string clientId)
        {
            return receivers.Contains(clientId);
        }
    }
}

编辑 2:

现在我的代码如下所示:

public void ReceiveNext(string clientId, string channel, Action<object> callback)
{
    var subscription = Disposable.Empty;
    subscription = messages
        .Where(message => message.Channels.Contains(channel))
        .Subscribe(message =>
                   {
                       if (message.TryDispatchTo(clientId, callback))
                           subscription.Dispose();
                   });
}

class MessageWrapper
{
    readonly object message;
    readonly List<string> receivers;

    public MessageWrapper(List<string> channels, object message)
    {
        this.message = message;
        receivers = new List<string>();
        Channels = channels;
    }

    public List<string> Channels { get; private set; }

    public bool TryDispatchTo(string clientId, Action<object> handler)
    {
        lock (receivers)
        {
            if (IsReceivedBy(clientId)) return false;
            handler(message);
            MarkAsReceivedFor(clientId);
            return true;
        }
    }

    void MarkAsReceivedFor(string clientId)
    {
        receivers.Add(clientId);
    }

    bool IsReceivedBy(string clientId)
    {
        return receivers.Contains(clientId);
    }
}

【问题讨论】:

  • 你为什么需要HandleNext?你不能在一个函数中处理整个队列吗? HandleNext 对我来说似乎不是很 rx'y (因为消息是可变的,而它不应该这样)。函数式编程通过不使用可变的东西来消除竞争。
  • 嗯...你是对的,我实际上是在轮询消息。我必须承认你让我想到了这里。我认为 rx 非常适合这个,现在我不太确定了。
  • 您的代码似乎有点做作。这是您正在使用的实际代码还是您在此处发布的“愚蠢”?你能解释一下如何拨打HandleNext 吗?另外,作为旁注 - Rx 可能非常适合您正在做的事情,但是您做这件事的方式正在扼杀它。
  • 好吧,我只是把它降低了一个档次 :) 它确实工作得很好。我只有这一项比赛条件测试失败。我现在根据您的评论更新了问题。
  • 好吧,我建议如下:处理代码只适用于整个队列。只要队列为空,代码就等待队列中的下一条消息。所以不需要标记已经处理的消息。语义将更像一个管道。你不需要显式调用HandleNext:消息一进入队列,worker就会被它推送。

标签: c# system.reactive


【解决方案1】:

在我看来,您正在为自己制造 Rx 噩梦。 Rx 应该提供一种非常简单的方法来将订阅者连接到您的消息。

我喜欢这样一个事实,即您有一个自包含的类来保存您的 ReplaySubject - 它可以阻止您代码中的其他地方恶意并过早调用 OnCompleted

但是,ReceiveNext 方法不提供任何方法让您删除订阅者。至少是内存泄漏。您在 MessageWrapper 中跟踪客户端 ID 也是潜在的内存泄漏。

我建议你尝试使用这种函数而不是ReceiveNext:

public IDisposable RegisterChannel(string channel, Action<object> callback)
{
    return messages
        .Where(message => message.Channels.Contains(channel))
        .Subscribe(message => callback(message.Message));
}

这是非常 Rx-ish。这是一个很好的纯查询,您可以轻松取消订阅。

由于Action&lt;object&gt; callback 无疑与clientId 直接相关,我会考虑在其中放置防止重复消息处理的逻辑。

现在你的代码是非常程序化的,不适合 Rx。似乎您还没有完全了解如何最好地使用 Rx。这是一个好的开始,但您需要更多地从功能上思考(如在函数式编程中)。


如果您必须按原样使用您的代码,我建议您进行一些更改。

Queue 中这样做:

public IDisposable ReceiveNext(
    string clientId, string channel, Action<object> callback)
{
    return
        messages
            .Where(message => message.Channels.Contains(channel))
            .Take(1)
            .Subscribe(message =>
                message.TryReceive(clientId, callback));
}

MessageWrapper 中去掉MarkAsReceivedForIsReceivedBy 并改为这样做:

    public bool TryReceive(string clientId, Action<object> callback)
    {
        lock (receivers)
        {
            if (!receivers.Contains(clientId))
            {
                callback(this.Message);
                receivers.Add(clientId);
                return true;
            }
            else
                return false;
        }
    }

虽然我真的不明白为什么你有 .Take(1),但这些更改可能会减少竞争条件,具体取决于其原因。

【讨论】:

  • 非常感谢您的详尽回答!至于你的最后一个建议,这正是我到目前为止所做的,但正如你所说 - 它不是 rx'y ......它使用锁。现在我对内存泄漏的事情有些怀疑:我的印象是,执行 .Take(1) 会在恰好一个元素之后自动取消订阅(这就是它存在的原因)。这实际上来自您自己的答案:stackoverflow.com/a/7707768/64105,也许我误解了?至于clientId的列表,我相信它们会被定期清除,因为ReplaySubject上有一个时间窗口?
  • 现在关于 RegisterChannel 方法,我不确定我是否明白。如果 Web 请求调用 RegisterChannel,然后收到消息,请求将响应然后关闭。现在,如果有另一条消息进来,则没有活动的请求来处理回调,消息就会丢失。 RegisterChannel 应该如何使用?它不会只是将取消订阅的责任转移到回调函数中吗?
  • @asgerhallas - 使用Take(1) 确实在一个值后取消订阅,但您现在期望在收到每个值后立即调用ReceiveNext。您的代码不会这样做。即使确实如此,但由于代码的多线程性质,您可能会错过消息。您真的想创建一个订阅一次并从中获取所有值的 Rx 查询。您永远都不想像现在一样继续重新订阅。
  • 但是谁会得到请求之间的值呢?我不明白在没有请求处理它们时调用的回调如何不会丢失?
  • @asgerhallas - 我认为这是您正在努力解决的问题。应该编写查询,以便每个新值都由一个请求处理 - 即 no Take(1)。相反,您应该从一个订阅中获得所需的每一个价值。
【解决方案2】:

我不确定像这样使用 Rx 是一个好习惯。 Rx 定义了流的概念,它要求没有并发通知。

也就是说,要回答您的问题,为了避免竞争条件,请在 IsReceivedByMarkAsReceivedFor 方法中加锁。

至于更好的方法,您可以放弃整个处理业务,在收到请求时使用ConcurrentQueueTryDequeue 消息(您只是在做Take(1) - 这适合队列模型)。 Rx 可以帮助您给每条消息一个 TTL 并将其从队列中删除,但您也可以在传入请求时这样做。

【讨论】:

  • 感谢您的回答。我宁愿不使用锁 - 我的印象是,那不是 rx 方式,我正在尝试学习它:) ...但如果我使用锁,那么放他们分别在这些方法中。我需要锁定整个 ReceiveNext 方法或@Enigmativity 在上面的答案中所做的那样,否则比赛仍然会发生。要了解您对队列的建议:如果我将其出列,则只有一个客户端可以处理它,并且我希望多个客户端能够处理它(如在 pub/sub 中)。还是我误会了你?
  • @asgerhallas - Rx 在后台使用了很多锁。不要害怕锁,只要小心它们。每当您使用锁时,您都在尝试解决并发问题,即使您认为自己拥有,这些问题也几乎不可能解决。关于锁定,请始终猜测自己。
  • @asgerhallas - 哦,在IsReceivedByMarkAsReceivedFor 方法中加锁将不会解决竞争条件。两者都必须在一个锁中,因此我的 TryReceive 方法。
  • @Enigmativity 是的,这不是我刚才在上面的评论中所说的吗? :) ...现在并不是我不知道如何适当地进行锁定,我只是认为 rx 会有与我的场景相匹配的运算符,所以如果我只是将它塑造成正确的路。我最喜欢在引擎盖下锁定。
  • @asgerhallas - 你没有以正确的方式使用 Rx,所以你不能指望你使用的锁定能正常工作。至少保持原子锁定。最好摆脱锁定。我的回答显示了一种无需锁定的方法。
猜你喜欢
  • 1970-01-01
  • 2011-04-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-06-12
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多