【问题标题】:Message queue in an IRC botIRC 机器人中的消息队列
【发布时间】:2012-07-28 01:31:25
【问题描述】:

我目前正在编写一个 IRC 机器人。我想避免过多的洪水,所以我决定创建一个消息队列,每 X 毫秒发送下一条消息,但我的 attempt 失败了。第 43 行:

unset.Add((string)de.Key);

抛出 OutOfMemory 异常。我完全不知道我做错了什么。

也许我还应该解释一下这种(可能很复杂)排队方式背后的总体思路。

首先,主要的Hashtable queueht 存储ConcurrentQueue<string> 类型,其中消息的目标用作键。我希望机器人遍历哈希表,从每个队列发送一条消息(如果队列被清空,则删除密钥)。我想不出一个合适的方法来处理哈希表本身,所以我决定创建另一个队列ConcurrentQueue<string> queue,它会在清空队列时存储键及其使用顺序。

假设一个队列中有数百个项目的假设情况(这是可能的),任何新请求都会被天知道延迟多长时间(消息之间的内置延迟加上延迟),所以我有方法 Add( ) 重建queue。我创建了queueht 的深层副本(或者我希望如此),并基于这个一次性副本生成一个新的queue,在此过程中将其删除。

我认为我的思路和/或代码非常错误,因为我几乎没有线程、比简单数组更复杂的集合和 OOP 习惯/约定的经验。我非常感谢通过解释解决我的问题。提前致谢!

编辑:发布整个课程。

class SendQueue
{
    Hashtable queueht;
    ConcurrentQueue<string> queue;
    Timer tim;
    IRCBot host;
    public SendQueue(IRCBot host)
    {
        this.host = host;
        this.tim = new Timer();
        this.tim.Elapsed += new ElapsedEventHandler(this.SendNewMsg);
        this.queueht = new Hashtable();
        this.queue = new ConcurrentQueue<string>();
    }
    public void Add(string target, string msg)
    {
        try
        {
            this.queueht.Add(target, new ConcurrentQueue<string>());
        }
        finally
        {
            ((ConcurrentQueue<string>)this.queueht[target]).Enqueue(msg);
        }
        Hashtable ht = new Hashtable(queueht);
        List<string> unset = new List<string>();
        while (ht.Count > 0)
        {
            foreach (DictionaryEntry de in ht)
            {
                ConcurrentQueue<string> cq = (ConcurrentQueue<string>)de.Value;
                string res;
                if (cq.TryDequeue(out res))
                    this.queue.Enqueue((string)de.Key);
                else
                    unset.Add((string)de.Key);
            }
        }
        if (unset.Count > 0)
            foreach (string item in unset)
                ht.Remove(item);
    }
    private void SendNewMsg(object sender, ElapsedEventArgs e)
    {
        string target;
        if (queue.TryDequeue(out target))
        {
            string message;
            if (((ConcurrentQueue<string>)queueht[target]).TryDequeue(out message))
                this.host.Say(target, message);
        }
    }
}

EDIT2:我知道while (ht.Count &gt; 0) 将被无限期执行。它只是以前版本的一部分,看起来像这样:

while (ht.Count > 0)
{
    foreach (DictionaryEntry de in ht)
    {
        ConcurrentQueue<string> cq = (ConcurrentQueue<string>)de.Value;
        string res;
        if (cq.TryDequeue(out res))
            this.queue.Enqueue((string)de.Key);
        else
            ht.Remove((string)de.Key);
    }
}

但是集合在评估时不能被修改(我发现很难),所以它不再那样了。我只是忘了更改while 的条件。

我冒昧地尝试了 TheThing 的解决方案。虽然它似乎实现了它的目的,但它并没有发送任何消息......这是它的最终形式:

class User
{
    public User(string username)
    {
        this.Username = username;
        this.RequestQueue = new Queue<string>();
    }
    public User(string username, string message)
        : this(username)
    {
        this.RequestQueue.Enqueue(message);
    }
    public string Username { get; set; }
    public Queue<string> RequestQueue { get; private set; }
}
class SendQueue
{
    Timer tim;
    IRCBot host;
    public bool shouldRun = false;
    public Dictionary<string, User> Users;  //Dictionary of users currently being processed
    public ConcurrentQueue<User> UserQueue; //List of order for which users should be processed
    public SendQueue(IRCBot launcher)
    {
        this.Users = new Dictionary<string, User>();
        this.UserQueue = new ConcurrentQueue<User>();
        this.tim = new Timer(WorkerTick, null, Timeout.Infinite, 450);
        this.host = launcher;
    }
    public void Add(string username, string request)
    {
        lock (this.UserQueue) //For threadsafety
        {
            if (this.Users.ContainsKey(username))
            {
                //The user is in the user list. That means he has previously sent request that are awaiting to be processed.
                //As such, we can safely add his new message at the end of HIS request list.

                this.Users[username].RequestQueue.Enqueue(request); //Add users new message at the end of the list
                return;
            }
            //User is not in the user list. Means it's his first request. Create him in the user list and add his message
            var user = new User(username, request);
            this.Users.Add(username, user); //Create the user and his message
            this.UserQueue.Enqueue(user); //Add the user to the last of the precessing users.
        }
    }
    public void WorkerTick(object sender)
    {
        if (shouldRun)
        {
            //This tick runs every 400ms and processes next message to be sent.
            lock (this.UserQueue) //For threadsafety
            {
                User user;
                if (this.UserQueue.TryDequeue(out user))            //Pop the next user to be processed.
                {
                    string message = user.RequestQueue.Dequeue();   //Pop his request
                    this.host.Say(user.Username, message);
                    if (user.RequestQueue.Count > 0)                //If user has more messages waiting to be processed
                    {
                        this.UserQueue.Enqueue(user);               //Add him at the end of the userqueue
                    }
                    else
                    {
                        this.Users.Remove(user.Username);           //User has no more messages, we can safely remove him from the user list
                    }
                }
            }
        }
    }
}

我尝试切换到ConcurrentQueue,它应该也可以工作(尽管以更线程安全的方式,并不是说我对线程安全一无所知)。我也尝试切换到System.Threading.Timer,但这也无济于事。我很久以前就没有想法了。

编辑:作为一个彻头彻尾的白痴,我没有设置 Timer 启动的时间。将 bool 部分更改为更改计时器的到期时间和间隔的 Start() 方法使其工作。问题解决了。

【问题讨论】:

  • 您需要添加一些代码,以便其他人能够给出您想要的答案。
  • 我确实添加了代码。它在GitHub link。
  • 实际上,我们最多期望一小段代码能够重现您的问题。在最坏的情况下,您在此处发布的部分代码显示了问题发生的位置。链接很容易死掉,代码的链接很容易发生巨大的变化。问题是为了整个社区的利益,因此即使在未来也需要作为一个整体来表示。 =)
  • 每当调用SendQueue.Add() 时都会出现问题。 Visual Studio 用OutOfMemoryException 将我指向前面提到的第43 行,这就是我所知道的。编辑:我明白了。到时候我会发布整个课程。

标签: c# hashtable bots irc concurrent-collections


【解决方案1】:

据我所知,您希望能够将用户按顺序排列以及他们的每个请求。

意思是,如果一个用户请求像 1000 个请求,其他人仍然可以发送他们的请求,并且机器人以 FIFO 方式为每个用户提供 1 个请求。

如果是这样,那么你需要的是一种方式,类似于这个功能:

class User
{
    public User(string username)
    {
        this.Username = username;
        this.RequestQueue = new Queue<string>();
    }

    public User(string username, string message)
        : this(username)
    {
        this.RequestQueue.Enqueue(message);
    }

    public string Username { get; set; }
    public Queue<string> RequestQueue { get; private set; }
}


///......................

public class MyClass
{
    public MyClass()
    {
        this.Users = new Dictionary<string, User>();
        this.UserQueue = new Queue<User>();
    }

    public Dictionary<string, User> Users; //Dictionary of users currently being processed
    public Queue<User> UserQueue; //List of order for which users should be processed

    public void OnMessageRecievedFromIrcChannel(string username, string request)
    {
        lock (this.UserQueue) //For threadsafety
        {
            if (this.Users.ContainsKey(username))
            {
                //The user is in the user list. That means he has previously sent request that are awaiting to be processed.
                //As such, we can safely add his new message at the end of HIS request list.

                this.Users[username].RequestQueue.Enqueue(request); //Add users new message at the end of the list
                return;
            }

            //User is not in the user list. Means it's his first request. Create him in the user list and add his message
            var user = new User(username, request);
            this.Users.Add(username, user); //Create the user and his message
            this.UserQueue.Enqueue(user); //Add the user to the last of the precessing users.
        }
    }

    //**********************************

    public void WorkerTick()
    {
        //This tick runs every 400ms and processes next message to be sent.
        lock (this.UserQueue) //For threadsafety
        {
            var user = this.UserQueue.Dequeue(); //Pop the next user to be processed.
            var message = user.RequestQueue.Dequeue(); //Pop his request

            /////PROCESSING MESSAGE GOES HERE

            if (user.RequestQueue.Count > 0) //If user has more messages waiting to be processed
            {
                this.UserQueue.Enqueue(user); //Add him at the end of the userqueue
            }
            else
            {
                this.Users.Remove(user.Username); //User has no more messages, we can safely remove him from the user list
            }
        }
    }
}

基本上,我们有一个用户队列。我们弹出下一个用户,处理他的第一个请求,如果他有更多请求等待处理,则将他添加到用户列表的末尾。

希望这能清除一些功能。为了记录,上面的代码更像是一个伪代码而不是功能代码xD

【讨论】:

  • 我已经使用了这个解决方案(只需添加一个Timer 和一个发送方法,同时使WorkerTick 成为计时器的新事件处理程序),但它不起作用。排队看起来不错,但它不发送任何东西(添加新用户队列或添加到现有用户队列时的断点已设置,与消息发送时的断点不同)。好吧,至少它不会发送任何异常。
  • 我不得不将 System.Timers.Timer 更改为 System.Threading.Timer,但除此之外,它就像一个魅力。
【解决方案2】:

据我所知,您永远不会从while 中逃脱,因为您永远不会从临时哈希表ht 中删除项目,直到它之外。因此,计数将始终为&gt; 0。

【讨论】:

  • 是的,当我不得不将 ht.Remove((string)de.Key) 搬到外面 WHILE 时,我错过了这一点。虽然这很可能是 OutOfMemory 异常的原因,但恐怕它不能解决任何问题。不过,感谢您指出这一点:)
  • 好的,那么您可能应该编辑您的问题以删除其中的那部分,因为这可能是某些人认为是您的问题的地方。 ;)
【解决方案3】:

试试这个:

class User
{
    public User(string username)
    {
        this.Username = username;
        this.RequestQueue = new Queue<string>();
    }

    private static readonly TimeSpan _minPostThreshold = new TimeSpan(0,0,5); //five seconds

    public void PostMessage(string message)
    {
        var lastMsgTime = _lastMessageTime;
        _lastMessageTime = DateTime.Now;
        if (lastMsgTime != default(DateTime))
        {
            if ((_lastMessageTime - lastMsgTime) < _minPostThreshold)
            {
                return;
            }
        }

        _requestQueue.Enqueue(message);     
    }

    public string NextMessage
    {
        get
        {
            if (!HasMessages)
            {
                return null;
            }

            return _requestQueue.Dequeue();
        }
    }

    public bool HasMessages
    {
        get{return _requestQueue.Count > 0;}
    }

    public string Username { get; set; }
    private Queue<string> _requestQueue { get; private set; }
    private DateTime _lastMessageTime;
}

class SendQueue
{
    Timer tim;
    IRCBot host;
    public bool shouldRun = false;
    public Dictionary<string, User> Users;  //Dictionary of users currently being processed
    private Queue<User> _postQueue = new Queue<User>();

    public SendQueue(IRCBot launcher)
    {
        this.Users = new Dictionary<string, User>();
        this.tim = new Timer(WorkerTick, null, Timeout.Infinite, 450);
        this.host = launcher;
    }

    public void Add(string username, string request)
    {
        User targetUser;
        lock (Users) //For threadsafety
        {
            if (!Users.TryGetValue(username, out targetUser))
            {
                //User is not in the user list. Means it's his first request. Create him in the user list and add his message
                targetUser = new User(username);
                Users.Add(username, targetUser); //Create the user and his message
            }

            targetUser.PostMessage(request);
        }

        lock(_postQueue)
        {
            _postQueue.Enqueue(targetUser);
        }
    }

    public void WorkerTick(object sender)
    {
        if (shouldRun)
        {
            User nextUser = null;

            lock(_postQueue)
            {
                if (_postQueue.Count > 0)
                {
                    nextUser = _PostQueue.Dequeue();
                }
            }

            if (nextUser != null)
            {                
                host.Say(nextUser.Username, nextUser.NextMessage);
            }
        }
    }
}

更新:在更好地理解需求后进行了更改。

这提供了每个用户的洪水控制和整体限制。它也简单得多。

请注意,这是即时编写的,甚至还没有编译,可能需要考虑一些围绕用户实例的线程问题,但它应该可以工作。

【讨论】:

    猜你喜欢
    • 2012-09-21
    • 1970-01-01
    • 2014-09-18
    • 2019-06-13
    • 1970-01-01
    • 2011-08-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多