【问题标题】:Rebus Sagas, Revisions and DeferredMessagesRebus Sagas、Revisions 和 DeferredMessages
【发布时间】:2018-05-21 12:33:56
【问题描述】:

我正在尝试配置一个按以下方式工作的 Saga:

  1. Saga 收到一条运输订单消息。该运输订单有一个 RouteId 属性,我可以使用它来关联同一“卡车”的运输订单
  2. 这些运输订单是由另一个系统创建的,该系统可以使用批处理来发送此订单。但是,此系统无法将同一地址的运输订单分组。
  3. 几秒钟后,我只用这个 RouteId 发送了另一条消息。我需要获取收到的 RouteId 的所有运输订单,按地址对它们进行分组,然后将其转换为另一个对象并发送到另一个 Web 服务。

但我面临两个问题:

  1. 如果我“同时”向第一个处理程序发送两条消息,每条消息都会出现,即使具有与该消息相关的属性,在处理第一条消息后 IsNew 属性也不会改变
  2. 在第二个处理程序中,我希望访问与这些 Saga 相关的所有数据,但我不能,因为数据似乎是那些消息的修订中的数据被延迟。

相关代码:

saga 的总线配置

Bus = Configure.With(Activator)
   .Transport(t => t.UseRabbitMq(rabbitMqConnectionString, inputQueueName))
   .Logging(l => l.ColoredConsole())
   .Routing(r => r.TypeBased().MapAssemblyOf<IEventContract(publisherQueue))
   .Sagas(s => {
       s.StoreInSqlServer(connectionString, "Sagas", "SagaIndex");
          if (enforceExclusiveAccess)
          {
              s.EnforceExclusiveAccess();
          }
       })
   .Options(o =>
       {
         if (maxDegreeOfParallelism > 0)
         {
            o.SetMaxParallelism(maxDegreeOfParallelism);
         }
         if (maxNumberOfWorkers > 0)
         {
            o.SetNumberOfWorkers(maxNumberOfWorkers);
         }
      })
   .Timeouts(t => { t.StoreInSqlServer(dcMessengerConnectionString, "Timeouts"); })
   .Start();

SagaData 类:

public class RouteListSagaData : ISagaData
{
    public Guid Id { get; set; }
    public int Revision { get; set; }

    private readonly IList<LisaShippingActivity> _shippingActivities = new List<LisaShippingActivity>();

    public long RoutePlanId { get; set; }

    public IEnumerable<LisaShippingActivity> ShippingActivities => _shippingActivities;
    public bool SentToLisa { get; set; }

    public void AddShippingActivity(LisaShippingActivity shippingActivity)
    {
        if (!_shippingActivities.Any(x => x.Equals(shippingActivity)))
        {
            _shippingActivities.Add(shippingActivity);
        }
    }

    public IEnumerable<LisaShippingActivity> GroupShippingActivitiesToLisaActivities() => LisaShippingActivity.GroupedByRouteIdAndAddress(ShippingActivities);
}

CorrelateMessages 方法

protected override void CorrelateMessages(ICorrelationConfig<RouteListSagaData> config)
{
    config.Correlate<ShippingOrder>(x => x.RoutePlanId, y => y.RoutePlanId);
    config.Correlate<VerifyRouteListIsComplete>(x => x.RoutePlanId, y => y.RoutePlanId);
}

如果 saga IsNew,则处理假定启动 Saga 并发送 DeferedMessage 的消息

public async Task Handle(ShippingOrder message)
{
  try
  {
    var lisaActivity = message.AsLisaShippingActivity(_commissionerUserName);

    if (Data.ShippingActivities.Contains(lisaActivity))
      return;

    Data.RoutePlanId = message.RoutePlanId;
    Data.AddShippingActivity(lisaActivity);
    var delay = TimeSpan.FromSeconds(_lisaDelayedMessageTime != 0 ? _lisaDelayedMessageTime : 60);

    if (IsNew)
    {
      await _serviceBus.DeferLocal(delay, new VerifyRouteListIsComplete(message.RoutePlanId), _environment);
    }
 }
 catch (Exception err)
 {
   Serilog.Log.Logger.Error(err, "[{SagaName}] - Error while executing Route List Saga", nameof(RouteListSaga));
   throw;
 }
}

最后是延迟消息的处理程序:

public Task Handle(VerifyRouteListIsComplete message)
{
  try
  {
    if (!Data.SentToLisa)
    {
      var lisaData = Data.GroupShippingActivitiesToLisaActivities();

      _lisaService.SyncRouteList(lisaData).Wait();

      Data.SentToLisa = true;
    }
    MarkAsComplete();
    return Task.CompletedTask;
  }
  catch (Exception err)
  {
    Serilog.Log.Error(err, "[{SagaName}] - Error sending message to LisaApp. RouteId: {RouteId}", nameof(RouteListSaga), message.RoutePlanId);
    _serviceBus.DeferLocal(TimeSpan.FromSeconds(5), message, _configuration.GetSection("AppSettings")["Environment"]).Wait();
    MarkAsUnchanged();
    return Task.CompletedTask;
  }
}

感谢任何帮助!

【问题讨论】:

    标签: rebus saga


    【解决方案1】:

    我不确定我是否正确理解了您正在经历的症状。

    如果我“同时”向第一个处理程序发送两条消息,每条消息都会出现,即使具有与该消息相关的属性,在处理第一条消息后 IsNew 属性也不会改变

    如果调用EnforceExclusiveAccess,我希望消息以串行方式处理,第一个使用IsNew == true,第二个使用IsNew == false。

    如果不是,我希望这两条消息与IsNew == true 并行处理,但是当插入 sage 数据时,我希望其中一条成功,另一条失败并返回 ConcurrencyException .

    在ConcurrencyException 之后,消息将被再次处理,这次是IsNew == false。

    这不是你正在经历的吗?

    在第二个处理程序中,我希望访问与这些 Saga 相关的所有数据,但我不能,因为数据似乎是那些消息的修订中的数据被延迟。

    您是说saga数据中的数据似乎处于VerifyRouteListIsComplete消息被延迟时的状态?

    这听起来很奇怪,也不太可能?你可以再试一次看看是否真的如此?


    更新:我找到了您遇到这种奇怪行为的原因:您不小心设置了您的 saga 处理程序实例以跨消息重复使用。

    你是这样注册的(警告:不要这样做!):

    _sagaHandler = new ShippingOrderSagaHandler(_subscriber);
    
    _subscriber.Subscribe<ShippingOrderMessage>(_sagaHandler);
    _subscriber.Subscribe<VerifyRoutePlanIsComplete>(_sagaHandler);
    

    Subscribe 方法随后在 BuiltinHandlerActivator 上进行此调用(警告:不要这样做!):

    activator.Register(() => handlerInstance);
    

    这是不好的原因(特别是对于 saga 处理程序),因为处理程序实例本身是有状态的——它有一个包含进程当前状态的 Data 属性,并且还包括 IsNew 属性.

    您应该始终做的是确保每次收到消息时都会创建一个新的处理程序实例——您的代码应该更改为如下内容:

    _subscriber.Subscribe<ShippingOrderMessage>(() => new ShippingOrderSagaHandler(_subscriber)).Wait();
    _subscriber.Subscribe<VerifyRoutePlanIsComplete>(() => new ShippingOrderSagaHandler(_subscriber)).Wait();
    

    如果Subscribe的实现改成这样可以做到:

    public async Task Subscribe<T>(Func<IHandleMessages<T>> getHandler)
    {
        _activator.Register((bus, context) => getHandler());
        await _activator.Bus.Subscribe<T>();
    }
    

    这将解决您的独占访问问题:)

    您的代码还有另一个问题:您在注册处理程序和启动订阅者总线实例之间存在潜在的竞争条件,因为理论上您可能会很不幸并在总线启动和注册处理程序之间开始接收消息。

    您应该更改代码以确保在启动总线之前注册所有处理程序(从而开始接收消息)。

    【讨论】:

    • 我在github.com/GersonDias/RebusSagaConcurrencyStackOverflow 上传了一个git repo,它重现了我的架构和我面临的错误。如果您运行 Rebus.Publisher 项目,它将向队列发布 4 条消息。之后运行 Rebus.Subscriber 项目,您将在数据库中看到除了延迟消息的代码之外还有 4 条延迟消息位于 if(IsNew) 内。我也调用了EnforceExclusiveAccess,并没有注意到并发异常。你能试着看看这段代码来帮助我找出我做错了什么吗?我希望和你一样
    • 但是,症状就是这样,消息不是以串行方式处理的,但是几乎正确地插入/更新了 saga(我不能确定,但​​我觉得有些事情搞砸了,尤其是集合和基于消息属性的正确关联)。
    • 问题应该出在我发送许多 IAmInitiatedBy 接口中定义的类型的消息吗?
    • 不,这不应该有所作为。如果您可以重现该问题,例如在单元测试或小型控制台应用程序中,我很乐意为您调试。
    • 非常感谢,@mookid8000!在这个 repo github.com/GersonDias/RebusSagaConcurrencyStackOverflow 中,您可以看到问题正在发生...我正在使用此代码进行的一项测试是启动 Rebus.Publisher 控制台应用程序以发送消息,然后它们启动 Rebus.Subscriber 项目。您将在数据库中看到 4 条延迟消息,除了调用发送消息之前的 if (IsNew)...
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-05-16
    • 2017-02-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-12-15
    • 1970-01-01
    相关资源
    最近更新 更多