【问题标题】:C# Handling N events asynchronouslyC# 异步处理 N 个事件
【发布时间】:2019-01-30 02:11:20
【问题描述】:

所以交易是我有一个处理程序正在侦听一个队列,消息从另一个服务推送到它。我知道有多少消息应该到达队列,但是如何在测试中验证这一点?

假设我有一个像这样的处理程序 seutp:

public class Handler<TMessage> : IHandleMessages<TMessage>
{
    public TMessage Message { get; private set; }

    public async Task Handle(TMessage message)
    {
        await Task.Run(() =>
            {
                Message = message;
            })
            .ConfigureAwait(continueOnCapturedContext: false);
    }
}

这样每当TMessage 被放入队列时,它就会被消耗掉。

现在我想要验证我确实收到了我期望的消息数量,我尝试了以下操作:

public async Task VerifyReceivedMessages()
{
    //I'm expecting 5 files (how and why is not relevant in this case)
    const int numberOfMessages = 5;
    int numberOfReceivedMessages = 0;

    var receivedMessages = new List<TMessage>();

    var handler = new Handler<TMessage>(new List<TMessage>());

    while (numberOfReceivedMessages < numberOfMessages)
    {
        handler.WaitForMessage();
        foreach (var message in handler.MessageList)
        {
            if (!receivedMessages.Contains(message))
            {
                receivedMessages.Add(message);
                numberOfReceivedMessages = receivedMessages.Count;
            }
        }
    }
    //I rarely get to this one before the program terminates.
    await SaveReceivedMessages(receivedMessages);
}

为此,我还扩展了处理程序,因此它现在包括一个等待消息的方法以及一个消息列表,它会在收到消息时添加到该列表中:

public class Handler<TMessage> : IHandleMessages<TMessage>
{
    public TMessage Message { get; private set; }
    public List<TMessage> MessageList { get; set; }

    public Handler(List<TMessage> messageList)
    {
        MessageList = messageList;
    }

    public async Task Handle(TMessage message)
    {
        await Task.Run(() =>
            {
                Message = message;
                ReceivedMessages(Message);
            })
            .ConfigureAwait(continueOnCapturedContext: false);
    }

    private void ReceivedMessages(TMessage Message)
    {
        MessageList.Add(message);
    }

    public void WaitForMessage(int sleepMsCycle = 100, int sleepLimitMs = 60000)
    {
        long maxWait = DateTime.UtcNow.Ticks + (sleepLimitMs * 10000);

        while(Message == null && DateTime.UtcNow.Ticks <= maxWait)
        {
            Thread.Sleep(sleepMsCycle);
        }

        long now = DateTime.UtcNow.Ticks;
        if (now > maxWait)
        {
            //Some exception handling...
        }
    }
}

处理程序本身似乎工作。每当发布服务发送消息时,处理程序都会将其拾取。问题是,发布者可以在不同的时间间隔发送 5 条消息,有时非常快,有时可能需要几秒钟。

就目前而言,我无法找到接收VerifyReceivedMessages() 消息的好方法。当前的解决方案有时只接收 1、2 或 3 个(在我看来,似乎是随机的),剩余的消息在队列中排队,然后在 foreach 的下一次迭代中终止。

有什么建议吗?

【问题讨论】:

标签: c# asynchronous message-queue


【解决方案1】:

所以我有机会睡在这个上面,结果发现它不起作用的原因是因为我试图从一个列表中读取,同时处理程序也可能在 @987654321 中访问该列表@。添加.ToList() 可以,但更好的是,将其放入线程安全列表也可以解决问题。

public ConcurrentDictionary<TMessage, string> MessageList { get; set; }

public Handler(ConcurrentDictionary<TMessage, string> messageList)
{
    MessageList = messageList;
}

...

private void ReceivedMessages(TMessage Message)
{
    MessageList.TryAdd(message, "");
}

【讨论】:

    猜你喜欢
    • 2011-09-11
    • 2011-08-03
    • 1970-01-01
    • 1970-01-01
    • 2017-07-28
    • 1970-01-01
    • 2021-04-16
    • 2012-08-23
    • 1970-01-01
    相关资源
    最近更新 更多