【问题标题】:NServiceBus events lost when published in separate threadNServiceBus 事件在单独的线程中发布时丢失
【发布时间】:2018-02-28 20:38:21
【问题描述】:

我一直致力于在 Azure 传输上使用 NServiceBus 获取长时间运行的消息。基于this document,我认为我可以在单独的线程中启动长进程,将事件处理程序任务标记为完成,然后监听自定义的 OperationStarted 或 OperationComplete 事件。我注意到大多数情况下我的处理程序都没有收到 OperationComplete 事件。事实上,只有在 OperationStarted 事件发布后我立即发布它时才收到它。两者之间的任何实际处理都会以某种方式阻止接收完成事件。这是我的代码:

用于长时间运行消息的抽象类

public abstract class LongRunningOperationHandler<TMessage> : IHandleMessages<TMessage> where TMessage : class
{
    protected ILog _logger => LogManager.GetLogger<LongRunningOperationHandler<TMessage>>();

    public Task Handle(TMessage message, IMessageHandlerContext context)
    {
        var opStarted = new OperationStarted
        {
            OperationID = Guid.NewGuid(),
            OperationType = typeof(TMessage).FullName
        };
        var errors = new List<string>();
        // Fire off the long running task in a separate thread
        Task.Run(() =>
            {
                try
                {
                    _logger.Info($"Operation Started: {JsonConvert.SerializeObject(opStarted)}");
                    context.Publish(opStarted);
                    ProcessMessage(message, context);
                }
                catch (Exception ex)
                {
                    errors.Add(ex.Message);
                }
                finally
                {
                    var opComplete = new OperationComplete
                    {
                        OperationType = typeof(TMessage).FullName,
                        OperationID = opStarted.OperationID,
                        Errors = errors
                    };

                    context.Publish(opComplete);

                    _logger.Info($"Operation Complete: {JsonConvert.SerializeObject(opComplete)}");
                }
            });

        return Task.CompletedTask;
    }

    protected abstract void ProcessMessage(TMessage message, IMessageHandlerContext context);
}

测试实施

public class TestLongRunningOpHandler : LongRunningOperationHandler<TestCommand>
{
    protected override void ProcessMessage(TestCommand message, IMessageHandlerContext context)
    {
        // If I remove this, or lessen it to something like 200 milliseconds, the 
        // OperationComplete event gets handled
        Thread.Sleep(1000);
    }
}

操作事件

public sealed class OperationComplete : IEvent
{
    public Guid OperationID { get; set; }
    public string OperationType { get; set; }
    public bool Success => !Errors?.Any() ?? true;
    public List<string> Errors { get; set; } = new List<string>();
    public DateTimeOffset CompletedOn { get; set; } = DateTimeOffset.UtcNow;
}

public sealed class OperationStarted : IEvent
{
    public Guid OperationID { get; set; }
    public string OperationType { get; set; }
    public DateTimeOffset StartedOn { get; set; } = DateTimeOffset.UtcNow;
}

处理程序

public class OperationHandler : IHandleMessages<OperationStarted>
, IHandleMessages<OperationComplete>
{
    static ILog logger = LogManager.GetLogger<OperationHandler>();

    public Task Handle(OperationStarted message, IMessageHandlerContext context)
    {
        return PrintJsonMessage(message);
    }

    public Task Handle(OperationComplete message, IMessageHandlerContext context)
    {
        // This is not hit if ProcessMessage takes too long
        return PrintJsonMessage(message);
    }

    private Task PrintJsonMessage<T>(T message) where T : class
    {
        var msgObj = new
        {
            Message = typeof(T).Name,
            Data = message
        };
        logger.Info(JsonConvert.SerializeObject(msgObj, Formatting.Indented));
        return Task.CompletedTask;
    }

}

我确定context.Publish() 调用正在被命中,因为_logger.Info() 调用正在向我的测试控制台打印消息。我还验证了它们被断点击中。在我的测试中,任何运行时间超过 500 毫秒的东西都会阻止 OperationComplete 事件的处理。

如果有人可以就在 ProcessMessage 实现中经过任何大量时间时 OperationComplete 事件未触发处理程序的原因提出建议,我将非常感激听到他们的声音。谢谢!

-- 更新-- 万一其他人遇到这个并对我最终做了什么感到好奇:

an exchange 与 NServiceBus 的开发人员合作之后,我决定使用实现 IHandleTimeouts 接口的 watchdog saga 来定期检查作业是否完成。我正在使用在作业完成时更新的 saga 数据来确定是否在超时处理程序中触发 OperationComplete 事件。这带来了另一个问题:当使用 In-Memory Persistence 时,saga 数据在线程间为not persisted,即使它被每个线程锁定。为了解决这个问题,我专门为长期运行的内存数据持久性创建了一个接口。该接口作为单例注入到 saga 中,因此用于跨线程读取/写入 saga 数据以进行长时间运行的操作。

我知道不建议使用 In-Memory Persistence,但根据我的需要配置另一种类型的持久性(如 Azure 表)是多余的;我只是想让OperationComplete 事件在正常情况下触发。如果在运行作业期间发生重新启动,我不需要保留传奇数据。无论如何,该作业将被缩短,如果作业运行时间超过设定的最长时间,saga 超时将处理触发 OperationComplete 事件并出现错误。

【问题讨论】:

  • 这个问题也得到了官方 NServiceBus GitHub 仓库的开发人员的回答:github.com/Particular/NServiceBus/issues/5121
  • 谢谢,@Sabacc,我更新了我的答案,以链接到我用 NServiceBus 打开的问题,并解释我的最终解决方案是什么。

标签: c# .net multithreading nservicebus


【解决方案1】:

原因是如果ProcessMessage足够快,你可能会在它失效之前得到当前的context,比如被销毁。

通过从Handle 成功返回,您是在告诉 NServiceBus:“我已经完成了这条消息”,因此它也可以对 context 执行它想要的操作,例如使其无效。在后台处理器中,您需要一个端点实例,而不是消息上下文。

当新任务开始运行时,您不知道Handle 是否已返回,因此您应该只考虑消息已被消费,因此无法恢复。如果您在单独的任务中发生错误,则无法重试。

避免没有持久性的长时间运行的进程。您提到的示例有一个存储来自消息的工作项的服务器,以及一个轮询此存储以查找工作项的进程。也许不理想,以防您横向扩展处理器,但它不会丢失消息。

为避免不断轮询,合并服务器和处理器,启动时无条件轮询一次,并在Handle 中安排轮询任务。注意这个任务只在没有其他轮询任务运行时才进行轮询,否则它可能会变得比持续轮询更糟糕。您可以使用信号量来控制它。

要横向扩展,您必须拥有更多服务器。对于某些 N,您需要衡量 N 个处理器轮询的成本是否大于以循环方式发送到 N 个服务器的成本,以了解哪种方法实际上执行得更好。在实践中,轮询对于低 N 就足够了。

为多个处理器修改示例可能需要较少的部署和配置工作,您只需添加或获取处理器,而添加或删除服务器需要更改指向它们的所有位置(例如配置文件)的端点。

另一种方法是将漫长的过程分解为多个步骤。 NServiceBus 有 sagas。这是一种通常针对已知或有限数量的步骤实施的方法。对于未知数量的步骤,它仍然是可行的,尽管有些人可能认为这是对 sagas 看似预期目的的滥用。

【讨论】:

  • 谢谢,@acelent。我最终使用了 NServiceBus 开发人员建议的看门狗传奇,但您的回答解决了上下文变得无效的原因。
猜你喜欢
  • 2013-09-23
  • 1970-01-01
  • 2015-07-18
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-08-17
  • 2011-10-11
相关资源
最近更新 更多