【发布时间】: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