【问题标题】:Implementing message tracking in MassTransit 3在 MassTransit 3 中实施消息跟踪
【发布时间】:2016-01-15 08:17:28
【问题描述】:

我已经将 MassTransit 与 Rabbit 队列一起使用了一段时间,并且已经实现了一个类似于此问题中的 IMessageTracker:How to log failed message in masstransit? 来跟踪消息重试,并在重试发生时为我们提供一些监控,以及在哪里消息的重试顺序 - 允许我们获取有关第三方端点何时导致我们重试的一些统计信息。这意味着我必须针对每个消息 ID 存储一个重试序列值,并在每次重试时检查它/增加它,然后在调用 MessageWasReceivedSuccessfully 时删除消息 ID。此外,如果消息在所有重试后都失败了,我们可以监控失败发生的时间(诚然,我们不需要消息跟踪器,但当我为重试创建跟踪器时,它变得很明智)。

我想升级到 MassTransit 3,因为新的重试策略在这种情况下会有所帮助,但我找不到使用 MassTransit 3 实现消息跟踪器的方法。

我目前的方法是为我的处理程序配置一个消费观察者,并为我的总线配置一个接收观察者,如下所示:

sbc.ReceiveEndpoint(host, config.QueueName, ep =>
{
    ep.Consumer(() => _container.Resolve<IEventConsumer>());
    ep.Observer(new EventConsumeObserver());
});        

_busControl.ConnectReceiveObserver(new ReceiveObserver(container.Resolve<IMonitoringPublisher>()));

观察者如下:

public class ReceiveObserver : IReceiveObserver
{
    public ReceiveObserver()
    {
    }

    public async Task PreReceive(ReceiveContext context)
    {
        await Console.Out.WriteLineAsync("ReceiveObserver - PreReceive Observed");
    }

    public async Task PostReceive(ReceiveContext context)
    {
        await Console.Out.WriteLineAsync("ReceiveObserver - PostReceive Observed");
    }

    public async Task PostConsume<T>(ConsumeContext<T> context, TimeSpan duration, string consumerType) where T : class
    {
        await Console.Out.WriteLineAsync("ReceiveObserver - PostConsume Observed");
    }

    public async Task ConsumeFault<T>(ConsumeContext<T> context, TimeSpan duration, string consumerType, Exception exception) where T : class
    {
        await Console.Out.WriteLineAsync($"ReceiveObserver - ConsumeFault Observed - {exception.Message}.");
    }

    public async Task ReceiveFault(ReceiveContext context, Exception exception)
    {
        await Console.Out.WriteLineAsync("ReceiveObserver - ReceiveFault Observed");
    }
}

public class EventConsumeObserver : IObserver<ConsumeContext<IEvent>>
{
    public void OnNext(ConsumeContext<IEvent> value)
    {
        Console.WriteLine("EventConsumeObserver - OnNext");
    }

    public void OnError(Exception error)
    {
        Console.WriteLine("EventConsumeObserver - OnError");
    }

    public void OnCompleted()
    {
        Console.WriteLine("EventConsumeObserver - OnCompleted");
    }
}

这确实允许我监视重试何时发生(有点麻烦),但我在消息跟踪器上找不到与 MessageWasReceivedSuccessfully 方法等效的方法。此外,当我运行一个模拟超时引发异常的测试时,输出似乎没有接近 on completed 方法 - 跟踪如下:

(模拟超时)

ReceiveObserver - PreReceive Observed
Request timeout
EventConsumeObserver - OnNext
ReceiveObserver - PostConsume Observed
Request timeout
EventConsumeObserver - OnNext
ReceiveObserver - PostConsume Observed
Request timeout
EventConsumeObserver - OnNext
ReceiveObserver - PostConsume Observed
ReceiveObserver - ConsumeFault Observed - Call timed out

(没有模拟超时)

ReceiveObserver - PreReceive Observed
Request success
EventConsumeObserver - OnNext

我希望在每次重试时都会有更明显的调用来消耗错误,或者在重试发生时调用 OnError 并在每个事件之后调用 OnComplete,但情况似乎并非如此。

我是否缺少某种类型的观察者,或者是否有其他方法可以挂钩重试,以便我可以监控它们何时发生以及在成功之前执行了多少次重试?

【问题讨论】:

  • ep.Observer() 正在创建一个新的消息使用者,所以我也会删除那个。
  • 即使删除了端点观察者,我仍然一无所获。在阅读了必须连接到端点的消费观察者之后,我将观察者连接到端点 - codegur.net/32913996/… - 这是我看到的相同行为,使用 MT 3.0.13
  • 哦,你可能会看到最近的版本,我认为观察者并没有在更旧的版本中工作。 3.1.2 是最新的。
  • 是的,这似乎使观察者工作!我的印象是 3.0.13 是最新的,我的错误。感谢您的帮助。
  • 这里和那里有一些未设置的连接已被修复。很高兴听到它现在有效!

标签: c# masstransit


【解决方案1】:

您可能想查看IConsumeObserver,它更符合实际消息消费而不是消息接收。虽然两者都携带有趣的信息,但应该为消费者抛出的每个异常调用消费者观察者。 Pre/Post 消费方法也会为每种消费的消息类型调用。

【讨论】:

  • 我也连接了一个消费观察者,使用 _busControl.ConnectConsumeObserver(new ConsumeObserver()),但是我从这里添加的代码中删除了它,因为我没有看到任何消息来自它(我有它只是在我运行测试时将消息记录到与其他观察者相同的控制台)。是不是我连接错了?
猜你喜欢
  • 2013-03-08
  • 1970-01-01
  • 2010-12-27
  • 1970-01-01
  • 2021-05-21
  • 1970-01-01
  • 2020-01-16
  • 1970-01-01
  • 2011-05-12
相关资源
最近更新 更多