【问题标题】:How to do guaranteed message delivery with SignalR?如何使用 SignalR 进行有保证的消息传递?
【发布时间】:2014-04-07 11:22:46
【问题描述】:

我正在使用 C# 和 SignalR 开发实时客户端-服务器应用程序。我需要尽快向客户端发送消息。 我在服务器上的代码:

for (int i = 0; i < totalRecords; i++)
{
    hubContext.Clients.Client(clientList[c].Key).addMessage(
    serverId, RecordsList[i].type + RecordsList[i].value);
    Thread.Sleep(50);       
}

如果延迟 >=50 毫秒,一切正常,但如果没有延迟或延迟小于 50 毫秒,则缺少一些消息。 我需要尽快发送消息,不要延迟。 我想我需要检查是否收到了消息,并且只有在发送另一条消息之后。
如何以正确的方式做到这一点?

【问题讨论】:

  • 您的服务器是在 IIS 中还是作为独立程序/服务运行?

标签: c# signalr


【解决方案1】:

扩展给出的答案,我做了以下事情:

我决定在发送消息的客户端使用 JS 中经过测试的 UUID 生成器之一为每条消息生成 UUID。

然后,将此 UUID 与消息一起发送。在其他客户端收到带有 UUID 的消息后,他将发送确认发送回发送者(确认包含所述 UUID)。

发件人收到他生成的消息 UUID 后,确定该消息已成功处理。

另外,在收到确认之前,我会阻止发送消息。

【讨论】:

    【解决方案2】:

    在收到来自其他客户端的确认之前,您需要重新发送消息。

    不要立即发送消息,而是将它们排队并让后台线程/计时器发送消息。

    这是一个可以工作的高性能队列。

    public class MessageQueue : IDisposable
    {
        private readonly ConcurrentQueue<Message> _messages = new ConcurrentQueue<Message>();
    
        public int InQueue => _messages.Count;
    
        public int SendInterval { get; }
    
        private readonly Timer _sendTimer;
        private readonly ISendMessage _messageSender;
    
        public MessageQueue(ISendMessage messageSender, uint sendInterval) {
            _messageSender = messageSender ?? throw new ArgumentNullException(nameof(messageSender));
            SendInterval = (int)sendInterval;
            _sendTimer = new Timer(timerTick, this, Timeout.Infinite, Timeout.Infinite);
        }
    
        public void Start() {
            _sendTimer.Change(SendInterval, Timeout.Infinite);
        }
    
        private readonly ConcurrentQueue<Guid> _recentlyReceived = new ConcurrentQueue<Guid>();
    
        public void ResponseReceived(Guid id) {
            if (_recentlyReceived.Contains(id)) return; // We've already received a reply for this message
    
            // Store current message locally
            var message = _currentSendingMessage;
    
            if (message == null || id != message.MessageId)
                throw new InvalidOperationException($"Received response {id}, but that message hasn't been sent.");
    
            // Unset to signify that the message has been successfully sent
            _currentSendingMessage = null;
    
            // We keep id's of recently received messages because it's possible to receive a reply
            // more than once, since we're sending the message more than once.
            _recentlyReceived.Enqueue(id);
    
            if(_recentlyReceived.Count > 100) {
                _recentlyReceived.TryDequeue(out var _);
            }
        }
    
        public void Enqueue(Message m) {
            _messages.Enqueue(m);
        }
    
        // We may access this variable from multiple threads, but there's no need to lock.
        // The worst thing that can happen is we send the message again after we've already
        // received a reply.
        private Message _currentSendingMessage;
    
        private void timerTick(object state) {
            try {
                var message = _currentSendingMessage;
    
                // Get next message to send
                if (message == null) {
                    _messages.TryDequeue(out message);
    
                    // Store so we don't have to peek the queue and conditionally dequeue
                    _currentSendingMessage = message;
                }
    
                if (message == null) return; // Nothing to send
    
                // Send Message
                _messageSender.Send(message);
            } finally {
                // Only start the timer again if we're done ticking.
                try {
                    _sendTimer.Change(SendInterval, Timeout.Infinite);
                } catch (ObjectDisposedException) {
    
                }
            }
        }
    
        public void Dispose() {
            _sendTimer.Dispose();
        }
    }
    
    public interface ISendMessage
    {
        void Send(Message message);
    }
    
    public class Message
    {
        public Guid MessageId { get; }
    
        public string MessageData { get; }
    
        public Message(string messageData) {
            MessageId = Guid.NewGuid();
            MessageData = messageData ?? throw new ArgumentNullException(nameof(messageData));
        }
    }
    

    这里是一些使用MessageQueue的示例代码

    public class Program
    {
        static void Main(string[] args) {
            try {
                const int TotalMessageCount = 1000;
    
                var messageSender = new SimulatedMessageSender();
    
                using (var messageQueue = new MessageQueue(messageSender, 10)) {
                    messageSender.Initialize(messageQueue);
    
                    for (var i = 0; i < TotalMessageCount; i++) {
                        messageQueue.Enqueue(new Message(i.ToString()));
                    }
    
                    var startTime = DateTime.Now;
    
                    Console.WriteLine("Starting message queue");
    
                    messageQueue.Start();
    
                    while (messageQueue.InQueue > 0) {
                        Thread.Yield(); // Want to use Thread.Sleep or Task.Delay in the real world.
                    }
    
                    var endTime = DateTime.Now;
    
                    var totalTime = endTime - startTime;
    
                    var messagesPerSecond = TotalMessageCount / totalTime.TotalSeconds;
    
                    Console.WriteLine($"Messages Per Second: {messagesPerSecond:#.##}");
                }
            } catch (Exception ex) {
                Console.Error.WriteLine($"Unhandled Exception: {ex}");
            }
    
            Console.WriteLine();
            Console.WriteLine("==== Done ====");
    
            Console.ReadLine();
        }
    }
    
    public class SimulatedMessageSender : ISendMessage
    {
        private MessageQueue _queue;
    
        public void Initialize(MessageQueue queue) {
            if (_queue != null) throw new InvalidOperationException("Already initialized.");
    
            _queue = queue ?? throw new ArgumentNullException(nameof(queue));
        }
    
        private static readonly Random _random = new Random();
    
        public void Send(Message message) {
            if (_queue == null) throw new InvalidOperationException("Not initialized");
    
            var chanceOfFailure = _random.Next(0, 20);
    
            // Drop 1 out of 20 messages
            // Most connections won't even be this bad.
            if (chanceOfFailure != 0) {
                _queue.ResponseReceived(message.MessageId);
            }
        }
    }
    

    【讨论】:

    • 考虑使用 BlockingCollection,参见msdn.microsoft.com/en-us/library/dd267312.aspx,无需手动锁定,效率更高。但它不能保证 FIFO 顺序。
    • 是的,这就是问题所在。就像我提到的,我没有花时间解决这个解决方案的锁定问题。 System.Collections.Concurrent 中可能有一个不错的对象,但这可能有点超出了这个范围。当然不是任何人都应该使用的最终解决方案,但它是一个很好的起点。
    • ConcurrentQueue 易于使用,可能正是您所需要的:devx.com/dotnet/working-with-concurrent-queue-in-c.html
    • 这已经快 3 年了,但 System.Threading.Channels 可能是您最好的选择。
    • @KellyElton 抱歉,我指的是关于队列的评论对话。
    【解决方案3】:

    SignalR 不保证消息传递。由于 SignalR 在您调用客户端方法时不会阻塞,因此您可以非常快速地调用客户端方法,正如您所发现的那样。不幸的是,客户端可能并不总是准备好在您发送消息后立即接收消息,因此 SignalR 必须缓冲消息。

    一般来说,SignalR 每个客户端最多可以缓冲 1000 条消息。一旦客户端落后超过 1000 条消息,它将开始丢失消息。这个 DefaultMessageBufferSize 的 1000 可以增加,但这会增加 SignalR 的内存使用,并且仍然不能保证消息传递。

    http://www.asp.net/signalr/overview/signalr-20/performance-and-scaling/signalr-performance#tuning

    如果您想保证消息传递,您必须自己确认它们。正如您所建议的,您只能在前一条消息被确认后发送一条消息。如果等待每条消息的 ACK 太慢,您也可以一次 ACK 多条消息。

    【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2015-05-17
    • 1970-01-01
    • 1970-01-01
    • 2018-09-05
    • 1970-01-01
    • 1970-01-01
    • 2023-04-06
    • 1970-01-01
    相关资源
    最近更新 更多